You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink KeyedStream的sum()函数是否有状态?状态支持checkpoint恢复吗?

我们先编写一个简单的词频统计(WordCount)作业:

DataStream<Tuple2<String, Integer>> counts =
        text.flatMap(new Tokenizer())
                .keyBy(value -> value.f0)
                .sum(1);

(数据源及其他无关细节不做讨论)
假设流入该处理管道的字符串为:

"the cat is on the table"

运行输出结果为:

<the - 1>
<cat - 1>
<is - 1>
<on - 1>
<the - 2>
<table - 1>

其中唯一出现两次的单词是"the"。
从运行现象来看,sum()函数似乎是有状态的:它至少会维护最新的<单词-计数>元组,当新的<word, 1>元组到达时更新对应值(按单词维度做键分区)。
如果该结论成立,那么在开启checkpoint的场景下,这部分状态是否会被存入检查点,并在任务发生故障时自动恢复?


解答

答案是肯定的,这部分状态默认会被Checkpoint持久化,故障时自动恢复。
你观察到的sum()维护状态的现象完全正确:keyBy之后调用的sum()是Flink内置的键控聚合算子,它内部为每个key维护了一个存储当前聚合值的键控状态,每来一条新数据就更新对应key的状态值、再向下游发送最新结果,和你自定义实现有状态聚合逻辑用的托管状态是同一套机制。

相关的关键细节:

  • 只要你在执行环境层面开启了Checkpoint(配置了env.enableCheckpointing(checkpointInterval)),所有内置算子的状态、你自定义的算子状态/键控状态,默认都会被纳入快照范围,不需要额外写代码把sum()的状态注册到Checkpoint,Flink框架自动完成这部分逻辑。
  • 故障恢复时,Flink会从最近一次成功完成的Checkpoint加载所有算子的状态快照,同时把可重放数据源(比如Kafka)的消费偏移量重置到快照对应位置,重放故障发生时未完成处理的数据,最终保证计数结果的一致性,不会出现计数清零、重复计数的问题。
  • 只有两种情况状态不会被持久化:一是手动给对应算子设置了setCheckpointingEnabled(false),二是自定义算子中没有通过Flink提供的StateDescriptor注册托管状态,而是用普通Java变量存储数据,内置的sum()算子不存在这两类问题。

内容的提问来源于stack exchange,提问作者Vin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.27 19:27:21