一、为什么“状态”是 Flink 的核心竞争力?

在实时计算里,单条事件往往无法给出答案:你需要累加、去重、会话识别、风控匹配、规则表关联……这都意味着跨事件记忆——也就是状态(State)。Flink 将状态与算子紧密绑定,并提供一致性保障(配合 checkpoint/savepoint),让你在大规模分布式场景下依然能写出确定性的实时应用。

二、Keyed DataStream:先按键分区,再谈状态

  • 只有在 KeyedStream 上才能使用 Keyed State
  • 使用 keyBy(KeySelector)(Python 为 key_by)将 DataStreamkey 重分区。
  • Key 是虚拟的(函数选择),并不要求把数据包装为键值对。

示例:用 POJO 的字段作为 Key

public class WC { public String word; public int count; public String getWord(){ return word; } }

DataStream<WC> words = // ...
KeyedStream<WC> keyed = words.keyBy(WC::getWord);

备注:Java 里的 tuple/expression key 方式已不推荐;KeySelector + Lambda 更直观、开销更小。

三、Keyed State 五大原语与用法

所有状态都以当前输入元素的 key 为作用域;可随时 clear() 清理当前 key 的状态。

  • ValueState:保存单值;update(T) / value()
  • ListState:保存列表;add / addAll / get() / update(List)
  • ReducingState:通过 ReduceFunction 将加入的元素规约为单值
  • AggregatingState<IN, OUT>:通过 AggregateFunction 聚合,输入类型可不同于聚合结果类型
  • MapState<UK, UV>:KV 映射;put / putAll / get / entries / keys / values / isEmpty

状态句柄获取三步走

  1. 定义 StateDescriptor(唯一名 + 类型信息 + 可选聚合函数)
  2. RichFunctionopen() 中通过 getRuntimeContext() 取状态句柄
  3. 在算子逻辑中读写状态

示例:两条聚一发的“简易计数窗口”

public class CountWindowAverage extends RichFlatMapFunction<Tuple2<Long, Long>, Tuple2<Long, Long>> {
  private transient ValueState<Tuple2<Long, Long>> sum; // f0: count, f1: sum

  @Override
  public void open(OpenContext ctx) {
    ValueStateDescriptor<Tuple2<Long, Long>> desc =
      new ValueStateDescriptor<>("average",
        TypeInformation.of(new TypeHint<Tuple2<Long, Long>>(){}),
        Tuple2.of(0L, 0L));
    sum = getRuntimeContext().getState(desc);
  }

  @Override
  public void flatMap(Tuple2<Long, Long> in, Collector<Tuple2<Long, Long>> out) throws Exception {
    Tuple2<Long, Long> cur = sum.value();
    cur.f0 += 1; cur.f1 += in.f1;
    sum.update(cur);
    if (cur.f0 >= 2) { out.collect(Tuple2.of(in.f0, cur.f1 / cur.f0)); sum.clear(); }
  }
}

四、状态 TTL:让状态“会过期”

为何需要 TTL?
控制状态大小、满足合规要求(如隐私数据在到期后不可读)、以及防止“僵尸 key”长期占用资源。

启用步骤

StateTtlConfig ttl = StateTtlConfig
  .newBuilder(Duration.ofMinutes(30))
  .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)      // 刷新时机
  .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) // 过期不可读
  .build();

ValueStateDescriptor<String> desc = new ValueStateDescriptor<>("v", String.class);
desc.enableTimeToLive(ttl);

关键选项

  • 刷新时机:

    • OnCreateAndWrite(默认)——创建/写入才刷新 TTL
    • OnReadAndWrite——读/写都会刷新(对某些场景更友好)
  • 可见性:

    • NeverReturnExpired(默认)——过期即视为不存在
    • ReturnExpiredIfNotCleanedUp——在后台尚未清理前,读可能返回过期值

注意事项

  • 仅支持处理时间 TTL;启用后会增加存储开销(Heap/RocksDB 均会存时间戳)。
  • 用开启 TTL 的描述符去恢复“之前未启用 TTL 的状态”(或反之)会触发兼容性错误。
  • 启用 TTL 后 StateDescriptordefaultValue 不再生效,请在业务层处理 null/过期 默认值。

五、过期状态的清理策略

思路:读时清理 + 后台清理。不同后端有不同手段,可灵活组合。

(一)读时清理
默认在调用 value() 等时对过期值做显式移除。

(二)全量快照时清理(Heap/RocksDB 全量快照)
cleanupFullSnapshot():在生成全量快照时做一次清理,缩小快照体积(从该快照恢复时不包含已删除的过期项)。

(三)增量清理(Heap)
cleanupIncrementally(checkedEntries, alsoPerRecord)

  • 每次访问状态都会推进“全局惰性迭代器”,检查并清理一定数量的条目;
  • 可选择在“每条记录处理”时也触发;
  • 会增加处理延迟;同步快照模式下会额外占内存(迭代器需要 key 副本)。

(四)RocksDB 压缩过滤清理
cleanupInRocksdbCompactFilter(queryEveryN, periodicCompactionTime)

  • 在 RocksDB compaction 期间用 TTL 过滤器判定过期条目并剔除;
  • queryEveryN 越小清理越快但 JNI 开销越大;
  • periodicCompactionTime 可加速“很少访问”的过期项清理(换取更多 compaction 次数)。

六、Operator State 与 Broadcast State:什么时候该用它们?

1)Operator State(非 keyed)

  • 绑定到并行算子实例,典型例子:Kafka Source 维护“分区→offset”映射。

  • 支持重分布(并行度变化):

    • 均分(even-split):把所有列表拼接后按并行度平均切分
    • 联合(union):每个实例恢复到同一整份列表(谨慎使用,高基数会让 checkpoint 元数据膨胀)

2)Broadcast State(特殊 Operator State)

  • 用于规则/维表广播:一个低吞吐“规则流”广播到所有下游实例,作为 map 状态被另一主数据流访问。
  • 同一算子可持有多个命名 broadcast states。

编程接口:CheckpointedFunction

public class BufferingSink implements SinkFunction<Tuple2<String,Integer>>, CheckpointedFunction {
  private transient ListState<Tuple2<String,Integer>> ckptState;
  private List<Tuple2<String,Integer>> buf = new ArrayList<>();
  private final int threshold;

  @Override
  public void snapshotState(FunctionSnapshotContext ctx) { ckptState.update(buf); }

  @Override
  public void initializeState(FunctionInitializationContext ctx) throws Exception {
    ListStateDescriptor<Tuple2<String,Integer>> d =
      new ListStateDescriptor<>("buffered-elements", TypeInformation.of(new TypeHint<Tuple2<String,Integer>>(){}));
    ckptState = ctx.getOperatorStateStore().getListState(d);
    if (ctx.isRestored()) for (Tuple2<String,Integer> e : ckptState.get()) buf.add(e);
  }

  @Override
  public void invoke(Tuple2<String,Integer> v, Context c) { /* 缓冲到阈值批量写出 */ }
}

七、有状态 Source:记得用 Checkpoint Lock

有状态 Source 需保证输出与状态更新的原子性(确保故障恢复下 exactly-once)。做法是使用 SourceContext#getCheckpointLock()

public static class CounterSource extends RichParallelSourceFunction<Long> implements CheckpointedFunction {
  private volatile boolean running = true;
  private Long offset = 0L;
  private ListState<Long> state;

  @Override
  public void run(SourceContext<Long> ctx) {
    final Object lock = ctx.getCheckpointLock();
    while (running) {
      synchronized (lock) { ctx.collect(offset); offset += 1; } // 输出与状态更新同锁保护
    }
  }

  @Override public void snapshotState(FunctionSnapshotContext c){ state.update(Collections.singletonList(offset)); }
  @Override public void initializeState(FunctionInitializationContext c) throws Exception {
    state = c.getOperatorStateStore().getListState(new ListStateDescriptor<>("state", LongSerializer.INSTANCE));
    for (Long v : state.get()) offset = v;
  }
}

如需在 checkpoint 完全确认后与外部系统交互,可实现 CheckpointListener

八、最佳实践清单

(一)设计阶段

  • 明确“状态作用域”:能 keyBy 就优先用 Keyed State;只有在无法按键分区时才考虑 Operator State

  • 按需选择状态原语

    • 单值:ValueState
    • 追加列表:ListState
    • 即时聚合:ReducingState/AggregatingState
    • 需要按键访存:MapState
  • 估算状态量级与增长速度,尽早规划 TTL清理策略

(二)编码规范

  • StateDescriptor 命名唯一、类型信息明确。
  • open()/initializeState() 里初始化状态句柄;只在算子逻辑中读写。
  • 读取 null/过期值时显式设置默认值,避免逻辑分支遗漏。

(三)TTL 与清理

  • 合规/隐私场景:NeverReturnExpired
  • 高频读场景:按需使用 OnReadAndWrite 刷新,平衡语义与性能。
  • Heap:优先启用增量清理;RocksDB:启用压缩过滤并按需调优 queryEveryN 与周期性 compaction。

(四)容错与扩缩容

  • 有状态 Source 要用 checkpoint lock
  • Operator State 选择合适的重分布模式(均分 vs 联合);
  • 大状态作业扩容前,尽量通过 savepoint 做平滑迁移。

(五)观测与排障

  • 建议监控:状态大小、TTL 清理速率、RocksDB compaction、反序列化失败、backpressure。

  • 常见问题:

    • 状态“越堆越大” → 缺失 TTL/清理策略;
    • 恢复失败 StateMigrationException → TTL 配置不兼容;
    • 规则广播不生效 → 检查是否正确使用 Broadcast State API 与流联接。

九、结语

掌握 Keyed State 原语 + TTL 清理 + Operator/Broadcast State 场景边界,基本就拿到了 Flink 有状态流处理的“驾驶执照”。接下来建议在回放环境用真实数据验证你的 TTL 与清理策略,跑通 checkpoint/savepoint 的全链路,再把作业推上生产。这样,你的状态既可靠可控

Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐