当一台设备持续上报测量值时,计算一分钟平均值似乎只是维护一个总和与一个计数器。但当设备数量增加,事件开始乱序,结果存储出现延迟,执行进程又随时可能退出时,问题就不再是“怎样求平均值”,而是:
当数据持续到来、事件发生乱序、下游处理变慢、执行节点发生故障时,Flink 如何组织分布式计算,并维持状态与结果的正确性?
本文围绕设备测量流展开,依次讨论执行模型、集群架构、状态管理、时间语义、一致性快照、背压、外部一致性、故障恢复与 SQL 计算。贯穿全文的主线是:计算不只是对当前事件执行一个函数,还要维护历史信息、判断时间进度,并为失败后的继续执行建立可靠边界。
阅读约定:本文以 Apache Flink 2.2 系列官方文档作为机制与 API 的说明基线,并不将其称为最新版本。主要讨论无界数据上的流式执行;批模式、状态后端和连接器的差异会单独说明。代码中标注为“片段”或“伪代码”的内容用于解释机制,不是完整部署工程。示例数据、并行度与吞吐量均为说明性设定,不是性能测试结论。全文配有 10 张独立技术图解;窗口生命周期、机制比较和结果变化等内容直接用表格说明。
一、Flink 解决什么问题:从消费消息到持续计算
1.1 一个持续统计程序,为什么需要计算引擎
假设设备上报的每条事件包含四个字段:
{ "event_id": "e-1001", "device_id": "device-A", "event_time": "2026-09-01T12:00:12.000Z", "measurement": 20.0}业务要求是:按照事件实际发生的分钟,持续统计每台设备的测量平均值。
最直接的实现,是启动一个消息消费者,将结果保存在进程内的 Map 中。下面的伪代码暂时不考虑乱序与输出时机:
// 伪代码:minuteOf 按事件时间计算分钟边界。Map<DeviceMinute, Accumulator> accumulators = new HashMap<>();
for (Event event : consumer.poll()) { DeviceMinute key = new DeviceMinute( event.deviceId(), minuteOf(event.eventTime())); Accumulator acc = accumulators.computeIfAbsent( key, ignored -> new Accumulator()); acc.sum += event.measurement(); acc.count += 1;}对于 device-A 在同一分钟内的测量值 20、22、24,程序维护 sum=66、count=3,平均值为 22。这部分计算没有分布式系统特有的难度。
困难来自程序的持续运行。
首先,累计结果不能随着进程退出一起消失。 如果程序处理了前两条事件后宕机,只恢复消费进度而没有恢复 sum=42、count=2,第三条事件就会成为新一轮统计的第一条,结果发生变化。
其次,重新启动必须知道从哪里继续读取。 对于分区日志,进度不是一个笼统的“已经消费”,而是每个分区各自的位置。多个分区同时参与计算时,还需要保留这些位置之间对应的状态边界。
最后,状态和读取进度必须匹配。 先保存状态、后保存位置,或者反过来,都可能在两次操作之间发生故障。只要它们不属于同一个可恢复边界,就可能多算或少算。
把 Map 改成数据库,只解决了状态放在哪里,并没有自动解决消费位置、计算状态与结果写入之间的一致性。增加多个消费者,又会引出同一设备的数据如何分配、状态由谁维护、扩容后怎样搬迁等问题。
Flink 将这些问题纳入运行时:它组织并行数据流,管理可恢复状态,并通过快照与输入重放,让计算能够跨越进程故障继续运行。窗口、时间进度与外部输出协议,则在这套基础上处理更具体的语义。12
因此,Flink 的职责不是替代一个 poll() 循环,而是管理持续计算的整个生命周期。
1.2 有界数据与无界数据
有界数据具有明确的结束位置。例如某个已经封存的文件集合,或者显式限定结束偏移量的一段日志。所有输入读取完毕后,系统可以知道不会再出现属于本次计算的新输入。
无界数据没有预先确定的终点。例如持续接收设备上报的主题。程序不能等待“所有设备事件都到齐”才输出结果,否则查询可能永远不产生结果。
二者的差异并不在于记录格式,而在于计算能够利用哪些事实:
| 维度 | 有界数据 | 无界数据 |
|---|---|---|
| 结束条件 | 可以确定输入已经结束 | 通常需要长期运行 |
| 结果表达 | 可以围绕完整输入产生最终结果 | 常按窗口输出,或持续更新已有结果 |
| 执行优化 | 可以利用完整输入、排序和物化中间结果 | 通常强调持续推进、增量计算和有限缓冲 |
| 故障恢复 | 可利用仍然有效的中间结果重新执行相关阶段 | 通常依赖状态快照与可重放输入 |
Flink 的 DataStream 支持 STREAMING 与 BATCH 执行模式。无界作业需要流式执行;有界作业可以采用批模式,也可以在流模式下运行。批模式能够使用不同的调度、数据交换、分组与聚合策略,在失败后利用可用的中间结果恢复,而不是照搬长期流作业的检查点路径。3
以同一批设备数据为例,流模式可能不断输出某台设备平均值的变化;批模式则可以在相应数据处理完成后给出最终平均值。对于适用的确定性计算,两者可以表达相同的最终关系结果,但到达结果的执行路径、输出次数和中间状态可能不同。
需要特别注意:依赖机器当前时间、外部可变查询或输入到达顺序的自定义逻辑,不能仅凭“批流统一”就假定执行结果完全相同。
批流统一指向计算表达与结果语义的统一,不意味着两种执行模式拥有完全相同的内部机制。
1.3 Flink 在数据系统中的位置
设备统计链路通常可以分为三个职责层:
数据源保留事实。 日志系统、文件系统或数据库变更流负责提供原始事件,并在各自能力范围内支持保留、读取和重放。
Flink 维护计算。 它解析、过滤、分区、关联与聚合事件,保存计算过程中需要的状态,并在故障后恢复这些状态。
结果系统提供访问。 聚合结果可以进入数据库、搜索系统、分析存储或新的消息主题,供查询、展示或后续计算使用。Source 与 Sink 连接器承担具体系统的接入工作。45
Flink 因而不是消息队列:它不能替代上游系统的日志保留职责。它也不是供任意业务请求长期查询的通用数据库:运行时状态首先服务于当前作业的计算与恢复,不等于一份对外提供通用事务和查询接口的业务表。
Flink 提供的不同 API,是表达计算的不同入口。DataStream 更适合显式描述事件处理、分区和状态操作;Table API 与 SQL 更适合以关系运算描述过滤、聚合与关联,再交由规划器转化为执行计划。声明 SQL 并没有绕开底层的数据流与状态,只是减少了手工组织这些细节的工作。67
设备统计链路可以概括为:
设备事件 → 可重放数据源 → Flink 持续计算 → 结果存储 → 查询应用托管状态服务于计算过程,检查点服务于故障恢复,结果存储服务于业务访问。这三种角色可以使用相似的存储技术,但不能因此视为同一份数据。
本章小结:Flink 解决的不只是读取数据,而是持续维护计算过程及其状态。
二、执行模型:一段程序怎样变成分布式数据流
2.1 Source、Operator 与 Sink
设备分钟均值可以表达为一条处理链路:
读取事件 → 解析字段 → 过滤无效数据 → 按设备分组 → 分钟窗口聚合 → 输出结果Source 是数据进入作业的入口。它负责与外部系统交互、获取记录,并通过连接器维护必要的读取进度。采用现代 Source API 时,通常还会区分发现与分配数据分片的 SplitEnumerator,以及真正读取分片的 SourceReader;它们的可恢复信息共同参与源端恢复。4
Operator 是处理数据的逻辑单元。解析、过滤、聚合、关联都可以形成算子。某些算子只依赖当前事件,某些算子需要访问历史状态,还可能响应定时器、Watermark 或快照请求。
Sink 则将结果送入外部系统。一个逻辑 Sink 不一定只有一个简单写入函数;需要提交语义的连接器,还可能展开写入者与提交者等运行时组件。5
下面是一段表达结构的伪代码:
source .map(parseEvent) .filter(isValidMeasurement) .keyBy(deviceId) .window(oneMinuteEventTimeWindow) .aggregate(sumAndCount) .sinkTo(resultSink);这里的每一步表示逻辑操作,并不表示启动一个独立进程。逻辑算子可能被展开为多个并行实例,也可能与相邻算子串联到同一个运行时任务中。理解分布式执行,需要进一步区分“算子定义”“算子并行实例”和“实际部署的 Task”。8
2.2 算子并行度与 Subtask
假设 Source 并行度为 2,解析和过滤的并行度为 2,窗口聚合的并行度为 4。一个窗口聚合算子就会展开为四个并行实例,通常称为四个 Subtask。
每个实例维护自己负责的数据与状态,而不是四个实例共同修改同一个全局累加器。设备 A 的当前窗口属于哪个实例,设备 B 又属于哪个实例,必须由数据分发规则确定。
对于 keyBy(deviceId),Flink 将相同 Key 的记录发送到同一个下游并行实例。这样,某台设备的累计值能够由一个明确的状态所有者维护。实际映射通过 Key Group 完成,第四章会展开。19
不同阶段可以使用不同并行度,是因为瓶颈可能不同。解析的成本、聚合的状态访问成本、外部结果系统的写入能力,不一定处于同一数量级。但是,并行度变化也意味着上下游需要建立对应的数据分发关系,而不是简单地复制几份代码。10
仍以上述设定为例:两个上游实例都可能接收到设备 A 的事件,但经过重新分区后,这些事件必须归入同一个聚合实例。因此,数据分发本身就是分布式计算的一部分。
这里有两个容易被忽略的限制。
其一,增加并行度不一定能分散单个热点 Key。如果设备 A 占据绝大多数流量,而业务仍要求所有 A 的状态由同一个 Key 管理,那么 A 仍会集中到一个实例。要拆分热点,需要重新设计可合并的局部聚合,而不只是修改并行度。
其二,按 Key 分组不等于按事件时间排序。来自不同分区或通道的记录可以乱序到达。keyBy 保证路由归属,不会自动把一个设备的所有事件按时间重新排好。
2.3 Operator Chaining 与网络数据交换
如果解析与过滤都是逐条处理,并行度一致,上下游采用兼容的一对一转发关系,而且满足链化策略及资源分组等条件,Flink 可以将它们串联成一个 Task。
在典型的链内处理过程中,解析算子处理一条记录后,直接调用下一个过滤算子的处理逻辑。同一个主任务线程可以连续执行多个算子,减少线程交接、缓冲以及序列化与反序列化开销。链化并不意味着一定不会发生对象复制;对象复用和类型相关行为仍有各自规则。86
假设一条设备 A 的事件进入了解析与过滤链:
同一个 Task 中:解析函数调用 → 过滤函数调用当它到达 keyBy 所定义的分区边界,就必须根据 Key 确定目标聚合实例。这不再只是“接着调用下一个函数”,而是进入跨任务的数据交换路径。
跨任务交换需要维护输出分区、数据缓冲和输入通道。两个 Task 位于不同 TaskManager 时,还涉及进程间网络传输;位于同一 TaskManager 时,可以采用本地交换路径,但它仍然不同于算子链内的直接调用。因此,发生 Shuffle 不代表每条记录一定跨越物理机器,关键在于发生了重新分区和任务边界交换。11
链化的收益也伴随约束。若链中某个用户函数长时间阻塞,整条链的主处理线程就不能及时处理下一条记录、定时器或部分控制工作。将算子放在同一 Slot,与将它们链化到同一个线程,也完全是两回事。

图 01:从逻辑算子到并行任务——链内直接调用、keyBy 分发与跨任务通道。
本章小结:追踪一条事件时,需要知道哪里只是函数调用,哪里改变了数据归属,哪里真正跨越任务或进程边界。
三、集群架构:谁负责协调,谁真正执行计算
3.1 JobManager 内部的职责划分
Flink 的运行时通常由 JobManager 与一个或多个 TaskManager 构成。前者协调作业,后者执行实际数据处理。用户事件在执行节点之间流动,并不是每条事件都先经过 JobManager 再转发。8
“JobManager”是协调侧的总体称呼,内部还需要区分三个职责:
| 组件 | 关注对象 | 主要职责 |
|---|---|---|
| Dispatcher | 作业提交与接入 | 提供作业提交入口,为作业启动对应的 JobMaster |
| ResourceManager | 集群计算资源 | 管理 TaskManager 提供的 Slot,并与部署环境的资源机制协作 |
| JobMaster | 单个作业 | 协调该作业的调度、执行状态、Checkpoint 与故障恢复 |
ResourceManager 解决“有哪些执行资源、怎样分配资源”。JobMaster 解决“这个作业的任务应当怎样在资源上运行,以及失败后怎样恢复”。二者协作,但不是同一层次的管理。8
例如,一个窗口 Task 失败后,JobMaster 需要确定受影响的恢复范围,并重新组织任务执行;可用 Slot 不足时,需要通过资源管理路径获取资源。至于能否启动新的 TaskManager,取决于部署方式。在 Standalone 环境中,Flink ResourceManager 本身并不会凭空创建新的工作进程;接入 Kubernetes 或 YARN 时,则可以利用相应环境的资源供应能力。12
这些组件名称表示逻辑职责,不能机械理解为必须运行三个互相独立的服务进程。具体进程布局与高可用方式由部署形态决定。
3.2 TaskManager、Task 与 Slot
TaskManager 是执行进程。它提供 CPU、内存、网络缓冲等运行资源,并承载多个 Task。典型 Java 部署中,它是一个 JVM 进程,也可能位于一个容器内。8
Task 是部署和运行的执行单元,可以包含一个算子的并行实例,也可以包含一条算子链。典型同步算子链由一个主任务线程推进;网络、异步 I/O、状态访问或快照工作还可能使用其他线程。
Slot 是资源分配与调度单位。它约束任务如何获得和共享 TaskManager 资源,但不是操作系统线程,也不是一枚绑定的 CPU 核心。给 TaskManager 配置四个 Slot,不意味着机器必须有四个核心,更不意味着每个 Slot 都独占一个核心。资源声明与调度预算,也不等于操作系统层面的 CPU 隔离。813
Slot Sharing 允许同一作业中属于同一共享组的不同阶段 Task,共享一个物理 Slot。假设 Source、Map 和窗口聚合的并行度都是 4,满足共享规则时,一个 Slot 可以容纳这三个阶段各自的一个执行实例。
这里要区分两种关系:
Operator Chaining:多个算子串联到同一个 Task 的执行链。Slot Sharing:多个 Task 共享调度资源,但仍可拥有各自的执行线程。因此,三个阶段各有四个实例,并不能直接得出“必须使用十二个物理 Slot”。在常见的单一 Slot Sharing Group、资源需求兼容的设定下,可以用四个共享 Slot 承载这些阶段。反过来,也不能无视资源需求和共享组,简单断言所有作业都只需要“最大并行度那么多 Slot”。13
共享资源提高了利用率,也引入共同故障和竞争边界。同一个 TaskManager 的多个 Task 会共享进程级开销;该进程退出时,它承载的任务也会一起受到影响。

图 02:Flink 集群架构——控制面、TaskManager、Slot 与 Task 的资源层级。
3.3 从提交作业到任务运行
将程序与集群对应起来,可以沿着五个连续阶段理解。
第一阶段是形成计算计划。 DataStream 程序构建转换关系;SQL 则经过解析、校验、关系计划优化与物理计划生成。流作业可以借助 StreamGraph、JobGraph 等层次表达计算拓扑、链化、并行度和数据交换方式。具体计划形成的位置与时机受执行模式和部署方式影响。613
第二阶段是作业接入与协调初始化。 提交路径把作业交给协调侧,由对应 JobMaster 管理执行图和执行尝试。
第三阶段是获取资源。 调度逻辑根据可运行任务与共享关系申请 Slot。是否需要额外创建 TaskManager,由现有资源及部署环境共同决定。
第四阶段是部署 Task。 代码、配置、输入输出关系,以及恢复时需要的状态句柄等信息被交给目标 TaskManager。工作进程初始化算子、连接数据通道,并在需要时加载状态。
第五阶段是开始处理。 Source 读取数据,记录沿算子链和跨任务通道推进。Checkpoint、心跳、失败通知等控制机制则在处理过程中持续参与协调。813
Session Cluster 与 Application Cluster 的区别,主要体现在生命周期和隔离边界。
Session Cluster 先启动共享集群,再接收多个作业。作业可以共享这组集群服务和工作资源,生命周期不必一致。它提高资源复用,但一个共享进程或资源层故障可能影响多个作业。
Application Cluster 围绕一个应用建立集群生命周期,应用入口 main() 在集群侧执行;一个应用仍然可能提交多个作业,因此不能把“一应用一集群”简单等同于“一 JobGraph 一集群”。应用之间的资源与故障隔离通常更明确,但仍依赖底层环境实际提供的隔离能力。12
本章小结:逻辑计算图决定要做什么,调度与资源体系决定这些计算在什么地方、以什么执行单元运行。
四、状态管理:Flink 怎样记住已经计算过的内容
4.1 普通变量与托管状态的区别
无状态算子可以只看当前记录。例如把 JSON 解析成事件对象,或者过滤缺少设备编号的数据。有状态算子则需要跨越多条事件保存信息:分钟均值需要累计值,去重需要已经见过的事件标识,关联需要等待另一侧数据的记录。1
普通成员变量也能“记住”信息,但 Flink 不会自动知道这些变量的业务含义。一个用户函数里的 HashMap 不会仅仅因为存在于 TaskManager 内,就自动获得快照、恢复与重新分配能力。
托管状态不同:程序通过 Flink 的状态接口访问数据,运行时知道状态属于哪个算子、哪个分区,并由状态后端提供存储与快照支持。应用当然也可以主动实现自定义状态快照接口,但那已经不是“普通变量会自动恢复”了。9
Keyed State 与 Operator State 是两种重要的归属方式。
Keyed State 按 Key 组织。 在 keyBy(deviceId) 之后,访问一个 ValueState,读到的是当前设备对应的值。处理设备 A 和设备 B 时,使用的可以是同一个状态句柄,但句柄背后定位到不同的状态条目。
Operator State 按算子并行实例组织。 它不以当前记录的 Key 为访问入口,适合表达某个实例所持有的分片信息或其他算子级数据。扩缩容时,需要依据状态所声明的分配规则重新分布,例如列表条目的重新划分;Union List State 与 Broadcast State 又有各自不同的语义,不能都当作普通列表均分。9
常见接口表达的是状态形状,而不是固定的存储介质:
| 接口 | 表达内容 | 设备场景中的用途 |
|---|---|---|
ValueState<T> | 当前 Key 的一个值 | 最近一次测量值,或一个累计器 |
MapState<K,V> | 当前 Key 内部的一组映射 | 事件 ID 去重表,或按分钟组织的累计器 |
ListState<T> | 一组条目 | 等待处理的数据;在 Operator State 中也可用于可重分配的条目集合 |
AggregatingState | 可增量更新的聚合累计器 | 保存求均值所需的总和与数量 |
下面的代码片段只展示一个设备累计器如何更新,不包含窗口划分和清理逻辑:
// 前提:已在初始化阶段注册名为 "device-accumulator" 的状态,// 并且该函数运行在 keyBy(deviceId) 之后。ValueState<Accumulator> accumulatorState;
void addMeasurement(double measurement) throws Exception { Accumulator previous = accumulatorState.value(); if (previous == null) { previous = new Accumulator(0.0, 0L); }
Accumulator next = new Accumulator( previous.sum + measurement, previous.count + 1);
accumulatorState.update(next);}状态句柄不是一个全局共享变量。对设备 A 的 update(),不会修改设备 B 的累计器。程序也应通过状态接口显式提交更新,不应依赖“修改取出的 Java 对象后,后端一定会感知”这一假设;对象是否直接对应堆内数据,取决于后端实现。914
对于一分钟一个累计器,可以让窗口算子管理窗口命名空间,也可以在自定义函数中按分钟维护 MapState,并自行实现触发与清理。两种方案表达的业务可能相近,但后者要求应用承担更多生命周期管理责任。
4.2 keyBy、Key Group 与状态归属
keyBy 首先是一条数据分发规则,不是聚合操作。
stream.keyBy(Event::getDeviceId);这行代码决定相同设备的数据应该去哪里,却没有规定这些数据要相加、求平均还是去重。真正更新累计值的是后续的聚合或状态处理算子。6
Flink 没有直接把每一个 Key 固定绑定到一个永久不变的 Subtask 编号。它在 Key 与执行实例之间引入了 Key Group,作为 Keyed State 重新分配的基本单位。一个有状态算子的 Key Group 数量对应它的最大并行度,而运行时的每个并行实例负责若干 Key Group。110
可以用以下概念关系理解:
业务 Key → 稳定的 Key Group 映射 → 当前并行度下的 Subtask 归属这里的哈希规则应由 Flink 实现维护,不应在外部系统随意仿写一个 hash % parallelism 并认为二者等价。
假设最大并行度为 12,即存在 KG0 到 KG11。当前并行度为 3 时,可以按连续范围理解为:
| 当前并行实例 | 负责的 Key Group |
|---|---|
| Subtask 0 | KG0~KG3 |
| Subtask 1 | KG4~KG7 |
| Subtask 2 | KG8~KG11 |
当并行度调整为 4,同一批 Key Group 可以重新分配为:
| 新并行实例 | 负责的 Key Group |
|---|---|
| Subtask 0 | KG0~KG2 |
| Subtask 1 | KG3~KG5 |
| Subtask 2 | KG6~KG8 |
| Subtask 3 | KG9~KG11 |
设备 A 如果原来属于 KG6,扩容后仍然属于 KG6,但承载它的实例从 Subtask 1 变为 Subtask 2。恢复系统按新的归属加载状态,数据分区规则也同步指向新的所有者。
最大并行度不是设备数量上限:一个 Key Group 内可以包含大量 Key。它限制的是这套状态划分可以支持的并行粒度。普通有状态恢复也不应随意修改最大并行度,因为那会改变状态分区布局,需要遵守相应的兼容与迁移机制。1015
还有一个边界:同一个 Key 在不同有状态算子里,各自拥有独立状态。设备 A 的去重状态与窗口累计状态,不会因为 Key 一样就自动合并成一份共享数据。

图 03:Key Group 与扩缩容——Key Group 保持稳定,承载状态的 Subtask 随并行度调整。
4.3 State Backend 与 Checkpoint Storage
理解状态存储,必须先拆开两个问题:
运行中的算子怎样读写状态?发生故障后,系统从哪里重新获得状态?
State Backend 回答第一个问题,也参与状态快照的实现。Checkpoint Storage 回答检查点数据与元数据如何保存、恢复时如何取得。二者可以共同使用文件系统,却不能因此被视为同一个概念。1416
对于常见的本地状态后端,正常处理时并不是每更新一次累计器,就把整个状态同步上传到对象存储。
HashMapStateBackend 主要将工作状态保存为 JVM 堆上的对象。访问路径短,但状态容量受到堆内存约束,且大对象集合会增加内存管理压力。形成检查点时,状态需要通过快照路径保存到配置的检查点存储。14
EmbeddedRocksDBStateBackend 使用嵌入式 RocksDB,将状态组织为字节形式的键值数据,工作集涉及内存缓存与本地磁盘文件。它可以承载超过 Java 堆容量的状态,但访问需要序列化、反序列化,并受磁盘、缓存和压缩合并等行为影响。
本地 RocksDB 文件不等于已经拥有跨节点容灾能力。TaskManager 所在机器或本地磁盘不可用时,恢复需要依靠可访问的持久化快照;可用的本地恢复副本只是加速路径,不能替代主恢复依据。1417
远程检查点存储常由分布式文件系统或对象存储承担。一个完成的检查点通常由元数据以及其引用的状态文件共同组成,并不保证“只有一个独立大文件”。增量检查点可以复用已有状态文件,只上传需要新增的部分;这里的增量通常体现为文件级复用,不是简单地把本轮修改过的业务 Key 列出来。1617
Flink 2.x 还提供 ForSt 等面向存算分离的状态能力。ForSt 可以把 SST 文件放到远程文件系统,本地磁盘用于缓存,并配合异步状态访问缓解远端访问延迟。因此,“Flink 的全部运行状态始终只在本地”并不是适用于所有后端的描述。本文采用的 2.2 文档仍将 ForSt 标为实验性能力,且它与传统后端在 State API、快照及 Savepoint 支持上存在差异,不能机械替换。1814
从这几条路径可以看出,选择后端是在平衡工作状态容量、读写延迟、序列化成本、快照开销以及恢复方式。把工作状态写得快,与把故障恢复点保存得可靠,是两个相关但不同的优化目标。

图 04:State Backend 与 Checkpoint Storage——运行时状态访问和持久化恢复是两条不同路径。
4.4 状态生命周期与清理
持续运行意味着持续积累。一个永不删除事件 ID 的去重集合,会随着不同事件数量增长;一个没有时间约束的流式关联,会持续保存仍可能与未来数据匹配的记录。
即使使用窗口,也不意味着状态天然很小。窗口长度、设备数量、乱序范围、允许迟到时间,以及是否保留原始记录,都会影响同时存活的状态规模。
对于均值,如果使用增量聚合,每个设备窗口只需要维护总和与计数器,空间开销主要与活跃的“设备与窗口组合”数量相关;若为了计算均值先保留全部原始测量值,状态大小则会随着窗口内事件数量增加。这个差异来自聚合表达,而不是 Checkpoint 配置。19
清理机制需要按职责区分。
窗口清理由窗口生命周期决定。窗口在达到清理边界后,释放其管理的内容和相关内部状态。自定义窗口函数额外注册的状态或其他用户资源,也应按接口要求清理,不能假设清理窗口等于清理程序中的一切。
定时器让应用在处理时间或事件时间到达某个条件时执行回调。例如,自定义每设备每分钟累计器,可以在事件时间边界到达时输出,再删除对应的 MapState 条目。事件时间定时器依靠 Watermark 推进,而不是单纯等待墙上时钟走过指定时刻。定时器也会参与故障恢复。20
State TTL为状态条目定义保留策略。在本文的 Flink 2.2 基线中,常规 State TTL 基于处理时间;可以配置写入时刷新,或读写时刷新,以及过期条目的可见性。过期条目在逻辑上不可见,并不意味着对应磁盘字节已经在同一时刻立即被物理删除,实际清理取决于后端与清理策略。9
TTL 也不是无限状态的万能开关。一个持续被访问并刷新 TTL 的 Key,可能长期不会过期;把去重记录保留十分钟,意味着十分钟之后到来的同一业务事件可能再次被接受。对持续聚合清理历史累计值,还可能把下一条输入当作新的统计起点。
状态清理实际上在回答:系统从什么时候开始不再记得某些历史?这个答案会直接影响业务结果。
本章小结:状态不只是一个 Map,还包含归属、访问、持久化、迁移与生命周期。
五、事件时间、水位线与窗口:乱序数据怎样参与计算
5.1 事件时间与处理时间
考虑设备 A 的三条事件:
| 事件 | 事件发生时间 | 到达计算节点的时间 | 测量值 |
|---|---|---|---|
| e1 | 12:00:12 | 12:00:13 | 20 |
| e2 | 12:00:50 | 12:00:51 | 22 |
| e3 | 12:00:40 | 12:01:04 | 24 |
事件 e3 实际发生在 e2 之前,却比 e2 更晚到达,并且跨过了机器时钟的分钟边界。
使用事件时间,三条事件都属于 [12:00:00, 12:01:00),这个窗口的正确输入均值为 22。使用处理时间,e3 则可能进入 [12:01:00, 12:02:00)。
这不是两种算法谁更准确的问题,而是统计口径不同。设备某一分钟真正发生了什么,应当以业务事件时间表达;系统当前一分钟处理了多少请求,则可能本来就应该以处理时间表达。21
事件时间把业务归属与处理速度解耦,但也引出新的问题:当机器时间已经进入下一分钟时,系统如何判断上一分钟还要不要等待?
仅仅存在 event_time 字段并不够。应用必须把它声明或提取为运行时可使用的时间戳,并配置时间进度策略。否则它只是一个普通字段,不会自动驱动事件时间窗口。2223
5.2 Watermark 怎样表达事件时间进度
Watermark 是数据流中的事件时间进度声明。可以将 Watermark(t) 理解为:该输入按照当前策略,声明自己在事件时间上已经推进到 t,早于这一进度的后续记录需要按相应的迟到规则处理。
它不是机器当前时间,也不是“所有更早事件已经被全世界确认到齐”的证明。设备离线补传、上游重试或异常延迟,都可能让更早的事件在水位线之后到达。2221
一种常用策略是有界乱序。设当前已观察到的最大事件时间为 Tmax,期望容忍的乱序范围为 Δ,则概念上可以理解为:
Watermark ≈ Tmax - Δ实际内置生成器还涉及毫秒边界处理与周期性发出,不能将这个近似公式当作每条事件都立即发布水位线的实现细节。
例如,Δ=5 秒,当某个输入观察到 12:01:08 的事件时,水位线大致可以推进到 12:01:03 附近。这并不表示记录先被固定缓存五秒,也不表示系统把所有记录重新排序;它是在“更早的事件还有多大概率继续到来”这一假设上决定时间进度。
下面是 Watermark 策略的 Java 片段,假设 eventTimeMillis 是业务事件时间的毫秒值:
WatermarkStrategy<Event> strategy = WatermarkStrategy .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, previousTimestamp) -> event.getEventTimeMillis()) .withIdleness(Duration.ofSeconds(30));对于支持分片级时间进度的 Source,通常应尽可能把策略放在 Source 上,让连接器利用分区或分片信息。若先将多个分区混合,再仅维护一个最大时间戳,可能失去对不同分区进度的细粒度感知。22
下游存在多个活跃输入通道时,通常以它们的最小 Watermark 作为可安全推进的共同进度:
输入 A 的 Watermark = 12:01:03输入 B 的 Watermark = 12:00:47下游共同进度 = 12:00:47原因很直接:B 仍可能带来属于上一分钟的记录,不能只因为 A 已经走得很远,就认定整个算子都已经走得很远。
这也解释了空闲输入为什么会阻塞窗口。一个长期没有事件的分区,其 Watermark 可能不再增长,从而压住整个算子的进度。空闲检测可以将满足条件的输入标记为 Idle,在计算活跃输入最小值时暂时排除它;该输入恢复后,重新到来的旧事件可能已经迟到。22
空闲检测不是“到时自动关闭所有窗口”的开关。当所有输入都空闲、没有新的事件时间进度时,不能假设水位线会凭空继续增长。对尾部窗口如何完成,需要结合有界输入结束、上游进度信号或明确的业务策略处理。

图 05:多输入 Watermark——下游受最慢活跃输入限制,空闲检测可暂时排除无数据分区。
5.3 窗口分配、触发与清理
窗口把持续到来的数据映射到有限的计算范围,但“分到哪个窗口”“什么时候输出”“什么时候删除状态”是不同的步骤。
滚动窗口按固定长度切分时间,相邻窗口不重叠。对于一分钟滚动窗口,12:00:40 属于 [12:00:00, 12:01:00),恰好发生在 12:01:00 的事件属于下一个窗口。
滑动窗口允许重叠。例如长度五分钟、每分钟滑动一次时,同一事件可能参与多个窗口。会话窗口则根据相邻事件之间的空闲间隔组织会话,后续事件可能连接原本分离的会话,因此涉及窗口合并。19
下面聚焦传统 DataStream 的事件时间滚动窗口,并假设使用默认事件时间触发器、允许迟到十秒。
窗口归属只由事件时间与窗口分配器决定。事件到达时,算子把它加入对应窗口,或在窗口已经无法接受数据时转入迟到处理路径。
首次输出由 Trigger 决定。默认事件时间触发器在 Watermark 达到窗口最大时间戳时触发。Flink 的时间窗口通常使用左闭右开区间,内部最大时间戳是 window_end - 1 毫秒。因此,“水位线推进到窗口结束附近”是方便理解的说法,精确实现还应考虑这一毫秒边界。19
允许迟到决定首次触发之后,窗口还保留多久,以便接收并处理迟到记录。它不会自动把首次输出再推迟十秒。对于该默认触发器,在清理边界前接受到的迟到事件可以导致再次输出更新后的窗口结果。
状态清理发生在清理时间条件满足时。此例中可以将其理解为窗口最大时间戳加上允许迟到时长。超过该边界的记录,不会让这个已经清理的窗口按原状态继续修正;可以按配置丢弃,或者进入侧输出进行另外处理。19
仍以三条测量值为例:
| 时刻或条件 | 窗口内部状态 | 对外结果 |
|---|---|---|
| 已收到 20 和 22,Watermark 未到触发点 | sum=42, count=2 | 尚未输出 |
| Watermark 达到首次触发条件 | 保留上述状态 | 输出平均值 21 |
| 清理前收到迟到的 24 | sum=66, count=3 | 再次输出平均值 22 |
| Watermark 达到清理条件 | 窗口状态释放 | 后续过迟事件走迟到处理路径 |
这里的 21 和 22 是同一设备同一分钟结果的不同版本,不是两个可以再次相加的独立测量值。窗口结果是否支持修改,还必须与后面的 Sink 语义配合。
以上描述针对明确配置的 DataStream 窗口。Flink SQL 的窗口聚合不能直接套用 DataStream 的 allowedLateness 接口及默认行为;SQL 的输出与迟到处理要依据对应查询算子的语义判断,第十章会单独说明。
5.4 Watermark 与 Checkpoint Barrier 的区别
Watermark 与 Checkpoint Barrier 都会在数据流中传播,但它们回答的是完全不同的问题。
Watermark 回答:业务事件时间推进到了哪里? 它影响事件时间定时器、窗口输出与时间相关状态清理。
Checkpoint Barrier 回答:哪些记录属于这次可恢复快照之前的计算范围? 它携带检查点身份,帮助不同任务建立一致的恢复边界。一个 Barrier 可以前后夹着完全乱序的事件时间,Checkpoint 也不要求数据先按时间排序。121
二者没有固定的一一对应关系:一个窗口存活期间可以完成多个 Checkpoint;一个 Checkpoint 也可能保存许多尚未结束的窗口。
“对齐”一词同样需要拆开。
Watermark Alignment协调多个输入的事件时间进度差异。当一个 Source 的事件时间进度远远领先其他 Source 时,可以暂停或放慢这个较快输入,避免下游为了等待慢输入而缓存过量数据。它限制的是“谁走得太快”,不是强行让慢输入宣布尚未达到的时间。该能力依赖现代 Source 及连接器提供相应支持。22
Checkpoint Barrier Alignment协调同一个检查点在多个输入通道上的边界。当某个通道的 Barrier 已到达,而其他通道还未到达时,需要防止边界之后的数据先进入当前算子的快照状态。它围绕 Checkpoint ID 工作,不围绕事件时间戳工作。
| 机制 | 协调对象 | 主要目的 |
|---|---|---|
| Watermark | 事件时间进度 | 决定时间驱动计算何时推进 |
| Watermark Alignment | 不同源的事件时间领先幅度 | 限制快输入导致的状态膨胀 |
| Checkpoint Barrier | 输入记录的快照边界 | 对应状态与输入处理进度 |
| Barrier Alignment | 多通道上的同一检查点边界 | 建立不混入边界后记录的局部快照 |
本章小结:数据什么时候发生、什么时候到达、什么时候输出,以及什么时候能够恢复,是不同维度的问题。
六、Checkpoint:持续运行时怎样形成一致性快照
6.1 为什么必须一起恢复状态与输入位置
先暂时忽略窗口关闭,只考察一个仍在累加的 sum。假设输入只有一个分区,并约定“处理到位置 100”表示第 100 条记录已经计入状态,恢复后下一条应读取 101。
在检查点 C42 的边界上,程序的事实是:
已经计入状态的输入:1 … 100sum = 1000恢复后的下一条输入:101之后,程序又处理了 101 至 120,这些记录贡献的总和为 200,因此当前 sum = 1200。
只比较位置 1 至 120 对累计值的贡献,三种恢复方式的差异如下:
| 恢复方式 | 恢复的 sum | 下一条读取位置 | 对原 1 至 120 条输入的累计结果 | 是否正确 |
|---|---|---|---|---|
| 使用 C42 中配套的状态与位置 | 1000 | 101 | 重放后为 1200 | 正确 |
| 较新状态搭配较旧位置 | 1200 | 101 | 重复计入 200,变为 1400 | 多算 |
| 较旧状态搭配较新位置 | 1000 | 121 | 跳过 101 至 120,仍为 1000 | 漏算 |
第一种方式允许代码重执行,但没有把同一段输入重复计入恢复后的状态。后两种方式的问题则不是“重启没有成功”,而是状态和输入进度描述了不同的历史边界。
正确恢复必须使用同一个检查点中的配套信息。也可以恢复一个更晚、已完成且自洽的快照,但不能自行混用不同边界的状态和输入位置。Checkpoint 的价值正是建立这种可恢复的一致性。2
真实作业通常有多个 Source Split 或分区,因此“输入位置”不是一个全局 offset,而是一组分区进度和相应的 Source 状态。算子状态、定时器以及参与协议的输出侧状态,也必须与这些进度对应。424
6.2 Barrier 如何建立快照边界
对一个正在运行的分布式作业来说,不能简单地广播一句“现在都把内存复制下来”。不同任务收到指令的时间不同,网络里还有尚未消费的记录,各个算子的状态也未必对应同一段输入。
Flink 使用流中的 Checkpoint Barrier 来协调这一边界。下面先讨论经典的对齐检查点,下一章再讨论它的非对齐变体。1
协调器触发一个带有唯一编号的检查点,例如 C42。Source 配合记录自身的读取进度,并在输出的数据流中放入 C42 的 Barrier。对于同一个有序通道,可以把记录理解成三段:
属于 C42 边界之前的记录 → Barrier(C42) → 属于 C42 边界之后的记录这里的“之前”和“之后”指流中的快照边界,不是记录携带的事件时间。12:00:50 的事件可以出现在 Barrier 前,12:00:20 的乱序事件也可以出现在 Barrier 后,这不会破坏 Checkpoint 协议。
对于只有一个输入通道的普通有状态任务,当它处理到 Barrier 时,此前的记录已经进入本地状态。任务就能够建立这一状态的快照视图,并继续向下游传播相同编号的 Barrier。不同任务会在不同的物理时刻完成自己的边界处理。1
因此,一致性快照建立的是数据流上的一致切面:恢复之后,不会出现某条记录已被计入下游状态、却既不属于上游已完成输出,也不在可重放或已保存的在途数据中的情况。
它不是“全世界同一毫秒的内存照片”,也不是把 TaskManager 中所有 Java 对象、线程栈与外部连接都打包。只有进入相应状态和协调协议的内容,才属于 Flink 能管理的恢复范围。242
6.3 多输入算子为什么需要对齐
假设一个聚合任务接收两个上游通道 A、B。对于 C42,它们的记录顺序如下:
通道 A:a1 → a2 → Barrier(C42) → a3 → a4通道 B:b1 → b2 → b3 → Barrier(C42) → b4A 的 Barrier 先到,而 B 的 b3 尚未处理完。
此时,任务不能直接保存快照:否则它的状态可能包含 A 边界前的全部贡献,却缺少 B 边界前的 b3。
它也不能一边等待 B,一边继续处理 A 的 a3、a4:这些记录属于 C42 之后,一旦进入当前状态,就把边界后数据混进了 C42 的快照。
对齐的做法是:暂时停止消费已经越过 C42 边界的 A 通道,继续处理尚未到达 C42 边界的 B 通道。 等到 B 的 Barrier 也到达,任务状态恰好包含 a1、a2、b1、b2、b3 的影响,而不包含 a3、a4、b4 的影响,此时再建立局部快照。117
这里“停止消费”不意味着立即让整个上游进程停止运行。A 后续的数据可能暂存在有限的网络缓冲中,缓冲耗尽后才进一步形成背压。其他尚未对齐的通道仍然可以被消费。
还有一个容易忽略的粒度:对齐面向输入通道,而不只是代码里的“有几个输入流”。 一个逻辑上只有一个输入的聚合算子,如果接收多个上游并行实例的数据,也可能拥有多个物理输入通道,因此同样涉及 Barrier 对齐。11
完成必要的局部快照动作后,任务可以继续处理记录;不需要等待整个作业所有任务都把快照文件上传完毕,才统一恢复运行。

图 06:多输入 Barrier 对齐——同一检查点边界到齐后形成一致的局部快照。
6.4 从局部快照到已完成 Checkpoint
从某个 Task 看见 Barrier,到 C42 成为可恢复边界,中间需要区分四个阶段:
| 阶段 | 已经完成什么 | 还不能据此认定什么 |
|---|---|---|
| 接收到 Barrier | 某个输入通道到达检查点边界 | 其他输入通道也已到齐 |
| 建立局部快照视图 | 固定本任务在该边界上的状态内容 | 所需状态文件已全部持久化 |
| 持久化并确认 | 本任务提交可恢复的状态句柄等信息 | 其他必要参与者也已完成 |
| 全局 Checkpoint 完成 | 协调器收齐必要确认,形成一致恢复点 | 任意外部副作用都已自动原子提交 |
对于支持异步快照的状态后端,可以把后续过程理解为:先在正确边界上形成稳定的快照视图,再异步写入或引用所需的持久化状态文件。前者需要与任务处理协调;后者尽可能与后续数据处理并行。异步并不表示保存一个持续被修改的普通对象引用,而是状态后端必须保证持久化读到的仍然是那个边界对应的内容。1416
完成本地快照持久化后,任务向协调器确认,提交用于恢复的状态句柄等信息。协调器等待该检查点要求的任务和协调组件完成确认,汇集元数据,才能将它标记为已完成。参与检查点提交协议的 Sink,则据此推进外部提交。245
因此,一个任务的快照文件已经出现在对象存储中,并不代表 C42 已经是有效的全局恢复点。若其他必要任务未完成,C42 仍可能超时或失败。故障恢复通常选择最新的、已经完成且仍可用的检查点,而不是挑每个任务“各自最新”的文件。16
检查点间隔也不能直接等同于最大重放时间。假设配置每 30 秒触发一次检查点,但快照完成较慢、受到并发数量限制、设置了最小间隔,或者连续几次失败,那么故障时能使用的最近完成边界可能已经明显早于 30 秒之前。17
更短的触发间隔有机会缩小需要重放的区间,但也增加快照协调和存储压力;如果这种压力反而拖慢作业与检查点完成,实际恢复收益可能下降。应同时观察完成频率、最近成功检查点的年龄、状态大小和恢复时的读取成本,而不是只看一个 interval。25
本章小结:Checkpoint 保存的是一组相互对应的状态和进度。其核心不是保存动作本身,而是这些分散的信息能够共同描述一个正确的计算边界。
七、背压与非对齐检查点:下游变慢后会发生什么
7.1 背压怎样沿数据流向上传播
假设结果存储发生延迟,Sink 的写入速度下降。Sink 接收的数据无法及时处理,输入缓冲逐渐被占满;上游任务无法继续向它输出,随后上游自己的输出缓冲也变得紧张。受限关系沿数据流反向传播,最终可能让 Source 降低读取速度。
数据本身仍然沿 Source 到 Sink 的方向流动;反向传播的是“下游暂时没有足够接收能力”这一约束。Flink 的网络缓冲与流量控制让生产速度受到消费能力约束,避免各个任务无限制地向内存堆积数据。11
可以用一个不改变记录数量的交换边界做估算:上游每秒产生 8 万条,下游每秒只能消费 5 万条,积压就以每秒 3 万条增长。即使额外缓冲能容纳 60 万条,也只能暂时吸收约 20 秒的差额。增加缓冲并没有把下游处理能力变成每秒 8 万条。
若 Source 来自可保留数据的日志系统,消费放慢后,更多积压会留在日志系统中;日志保留期限仍然是独立约束。背压并不意味着 Flink 可以保证任意上游都无限保存未读数据。
更重要的是,背压是一种现象,不是唯一的故障诊断。真正瓶颈可能是某个算子的 CPU、热点 Key、状态访问、磁盘、网络,或者外部系统的吞吐限制。最慢的任务可能表现为持续忙碌,而其上游任务才表现为明显背压。26
7.2 为什么背压会拖慢 Checkpoint
在对齐检查点中,Barrier 沿数据通道传播。如果前面已有大量尚未处理的记录,Barrier 就要等待这些记录向前推进;多个输入的积压程度不同,又会进一步拉大同一个检查点的 Barrier 到达时间差。
于是出现一条相互关联的因果链:
下游处理变慢 → 数据在有限缓冲中积压 → Barrier 到达下游更晚 → 多通道边界到达时间差扩大 → 对齐和检查点完成被拖慢作业仍然在处理记录,并不意味着 Barrier 已经有机会到达所有关键位置。在持续背压下,系统甚至可能一直“有吞吐”,却很久没有产生新的成功检查点。27
分析时,应区分几个指标:
| 指标 | 主要观察的问题 |
|---|---|
| Checkpoint Start Delay | 从检查点触发到该任务开始处理其第一个 Barrier,等待了多久 |
| Alignment Duration | 多个输入通道的检查点边界到齐需要多久 |
| Synchronous Duration | 建立局部快照等同步工作占用任务执行的时间 |
| Asynchronous Duration | 异步快照阶段持续多久,例如状态上传、相关 I/O 与等待 |
| End-to-End Duration | 从全局触发到检查点完成的总时间 |
这些指标描述的是不同阶段,不能把不同任务的各项耗时随意相加。尤其在非对齐检查点中,异步阶段还可能包含等待剩余通道边界和保存通道数据的时间,不能看到 Async Duration 高就立即断言“对象存储写入慢”。25
如果主要问题是第一个 Barrier 到达太晚,关注缓冲积压与上游推进;如果是对齐耗时很高,关注各通道速度差异;如果边界建立很快但持久化很慢,再重点检查状态量、存储带宽与快照方式。多个问题也可能同时存在。
7.3 Unaligned Checkpoint 怎样改变快照内容
对齐检查点的思路是:等相关输入都到达同一个边界,再把算子状态保存下来。
非对齐检查点换了一种做法:不把所有等待都放在“先消化缓冲、再对齐”这一步,而是把建立一致性所需的部分在途数据也纳入快照。 这样,Barrier 可以通过优先处理机制越过某些排队的数据,减少对缓冲排空速度的依赖。27
此时,一个可恢复的检查点不能只包含聚合器的 sum、count。它还需要描述相应通道中尚未被纳入算子状态、但恢复时必须继续处理的数据,也就是 Channel State。
可以从恢复需求反推其正确性。假设一条记录已经离开 Source 的检查点边界,但故障时尚未计入下游状态。若只恢复 Source 进度和下游状态,这条记录可能无处可找。把必要的在途数据保存下来,恢复时重建它在通道中的待处理状态,就补齐了这个缺口。
这不是“数据爱重复就重复”,也不是把 Exactly-once 降级为 At-least-once。它仍然要求每条记录在恢复后的状态演进中具有正确的归属,只是快照从“算子状态与输入进度”扩展到包含相应的通道状态。127
恢复时,Flink 协调还原算子状态、通道状态与 Source 进度,使保存的在途记录继续参与处理。不能自行再从一个更早的任意位置重放所有输入,否则同样可能制造重复。
还要区分“在途数据”与“所有积压”:非对齐检查点并不是把 Kafka 中尚未读取的全部历史消息复制进检查点,它保存的是该次一致性切面需要的相关网络缓冲数据。

图 07:背压下的对齐与非对齐检查点——后者把必要的在途数据纳入恢复边界。
7.4 非对齐检查点的收益与限制
非对齐检查点最直接的收益,是降低检查点对缓冲排空和多通道等待的敏感程度,使背压中的作业更容易持续产生可用恢复点。它改善的是检查点机制在拥堵条件下的推进方式,不是瓶颈算子的业务吞吐能力。27
代价也很明确:保存通道数据增加快照体积和存储 I/O,恢复时还要读取并重新处理这些数据。若瓶颈本来就在检查点存储端,增加写入量可能得不偿失。使用缓冲去膨胀等机制减少在途数据,可以同时影响检查点耗时和通道状态规模,但也需要结合实际吞吐观察。2717
另一个限制来自任务执行线程。如果一条记录的用户逻辑长时间阻塞,或者一次定时器触发执行大量同步工作,任务无法及时回到可处理检查点控制事件的位置。非对齐并不意味着运行时可以在任意一条 Java 语句中间强行冻结并持久化完整业务执行栈。27
版本与执行图的支持条件也不能忽略。以本文采用的 Flink 2.2 文档为基准,非对齐检查点面向 Exactly-once 模式,并限制检查点并发等配置;具体算子与连接方式还应满足对应支持条件。不要把“开启一个开关”理解为任意拓扑、任意状态后端和任意配置组合都具有相同效果。27
从机制上看,对齐与非对齐的取舍可以概括为:前者更多依靠等待把边界整理清楚,后者允许把尚未消化的必要数据一起带入快照。前者可能付出等待成本,后者可能付出存储与恢复成本。
本章小结:背压、网络缓冲与 Checkpoint 构成同一个运行系统。调整其中任何一项,都可能改变另外两项的表现,但不会凭空消除处理能力不足。
八、端到端 Exactly-once:引擎保证在哪里结束
8.1 状态一致性不等于函数只执行一次
继续使用第六章的例子:从 C42 恢复后,101 至 120 必须重新执行。解析函数、过滤条件、聚合逻辑都可能再次运行。如果在函数里打印日志,同一条输入甚至可能留下两次日志。
这并不自动否定 Flink 的 Exactly-once 状态语义。因为故障前但未进入恢复边界的状态修改也被回退了;重放是在恢复后的旧状态上继续演进,而不是在故障前已经累加完成的最新状态上再加一次。这里保证的是输入对受管理计算状态的影响,不是物理 CPU 指令只执行一次。2
这个保证还必须与上游业务重复区分。假设同一个 event_id = e-1001 因生产者重试,分别出现在日志的两个不同位置。对 Source 来说,这是两条输入记录。若计算逻辑没有按 event_id 去重,Checkpoint 不会自动判断它们其实代表同一次设备测量。
因此,至少要分清三件事:代码因恢复而重执行、同一输入对托管状态的重复影响,以及多条输入记录代表同一个业务事件。只有明确了“同一次”的身份,才能定义业务去重。
此外,状态一致性也不等于任意程序都会产生完全相同的重放值。若用户逻辑依赖处理时刻、随机数,或者随时变化的外部查询结果,即使读取了相同输入,也可能计算出不同内容。需要可重现结果时,应另外约束这些非确定性来源。28
8.2 外部写入为什么形成新的故障窗口
考虑如下过程:
处理输入 e-1001 → 外部数据库写入成功 → 对应的全局 Checkpoint 尚未完成 → 任务失败 → 从上一个已完成 Checkpoint 恢复 → e-1001 被重放并再次写入Flink 可以回退自己的托管状态,但不会自动撤销一次普通 JDBC 写入、一次 HTTP 请求,或者一次发送短信的调用。外部效果是否重复,取决于连接器协议与外部操作本身,而不能只看引擎启用了哪种检查点模式。29
不同写入动作面对重试有不同含义:
| 外部动作 | 重试后的典型风险 | 需要进一步明确的约束 |
|---|---|---|
| 追加一条没有稳定业务键的记录 | 同一结果出现两行 | 是否有可复用的结果身份或事务提交 |
| 按稳定主键覆盖完整结果 | 相同值重写通常不增加行数,但旧值可能覆盖新值 | 更新顺序、版本和并发写入条件 |
执行 total = total + delta | 同一增量可能再次累加 | 增量身份与原子去重 |
| 调用会产生现实副作用的接口 | 可能重复通知、重复创建业务对象 | 接口是否支持幂等键与结果查询 |
这张表不是某个数据库连接器的统一保证,而是根据操作语义进行的故障分析。例如“覆盖”在同一值被重复写入时可能是幂等的,但不能由此推出所有覆盖场景都是正确的。
8.3 Source、Checkpoint 与 Sink 如何协同
完整的端到端保证需要共同建立三个边界:Source 能够回到指定输入进度;Flink 能够恢复对应计算状态;Sink 能够把外部可见结果与该进度协调起来。缺少任何一方,端到端语义都需要重新判断。29
以可重放日志源和事务型 Sink 为例,可以把一次检查点关联的写入简化为事务 T42。实际实现可能是多个 Writer 分别维护多个事务,这里的单事务只是协议示意。
正常处理阶段:写入,但尚未对目标读者公开。 Writer 把该阶段产生的结果写入事务。外部系统能够接收这些数据,但应以适当的隔离语义屏蔽尚未提交的内容。
检查点阶段:准备提交并保存提交所需信息。 Writer 刷出待写数据,形成可恢复的提交信息,让它与检查点状态一起进入协调过程。现代 Flink Sink API 将预提交工作与最终提交工作分开,由 CommittingSinkWriter 产生 committable,再由 Committer 完成提交;具体事务组织由连接器实现。5
全局检查点完成之后:推进正式提交。 一旦 C42 已成为可恢复边界,Sink 的提交协议才允许对应结果进入正式可见状态。不是某个 Writer 自己保存完状态,就立即认为整个输入与计算过程已经安全完成。
现在重新分析故障窗口。
如果 C42 尚未完成就失败,恢复到更早的检查点时,未完成区间的事务应按连接器协议被中止或清理,输入重新计算后进入新的提交周期。
如果 C42 已完成,但任务还没有执行最终提交就失败,恢复路径需要找回待提交信息,继续完成这次提交。若提交已经成功,只是确认响应丢失,重复执行提交也必须能够识别事务身份,而不能产生第二份可见结果。
因此,外部系统需要支持对应的事务恢复、提交识别与超时管理,连接器也必须正确接入这些能力;“有事务”三个字本身并不足以完成这套协议。305
Kafka 是一个具体例子:使用支持该语义的 KafkaSink 时,Exactly-once 需要启用检查点,并使用不会与其他作业冲突的事务身份前缀;事务超时要覆盖正常检查点和可能的恢复过程;下游消费者需要采用 read_committed,才能按照已提交事务观察结果。30
输入端还存在一个常见误区:Kafka 消费者组已经提交的 offset,并不是 Flink 从检查点恢复时唯一或主要的进度依据。支持检查点的 Kafka Source 使用被纳入 Flink 状态的读取进度进行恢复;向 Kafka 提交 offset 还承担进度观测等作用,不能把二者混为一个原子操作。30
最后,这套机制不应被扩张为“所有外部系统在同一瞬间一起提交”。一个作业同时写 Kafka、数据库和另一个 HTTP 服务,不会因此自动获得跨三套系统的全局原子事务。保证必须分别落实到实际参与协议的输出边界。

图 08:端到端 Exactly-once——输入进度、算子状态、检查点完成与事务提交共同构成一致性边界。
8.4 幂等输出与业务约束
并不是所有结果存储都适合长事务。有些系统更适合使用幂等写入构造可接受的外部结果语义,但必须先定义稳定的写入身份。
对分钟平均值来说,结果身份可以是:
result_key = (device_id, window_start, window_end)payload = (sum, count, average, revision)同一个窗口的完整结果写到同一行,至少能避免“重试一次就多一行”。但允许迟到更新后,同一结果键还可能有多个有效版本。
假设版本 10 的结果是 21,随后版本 11 更新为 22。如果一次较早请求因为网络延迟,在版本 11 之后才到达数据库,普通 upsert 仍可能把 22 覆盖回 21。稳定主键解决的是身份,版本检查解决的是新旧顺序,这不是同一个问题。
可以让结果存储仅接受更高版本,或采用带预期版本的条件更新。版本必须有明确来源,在重试、重放和多写入者条件下仍然保持可解释的顺序;不能每次重试都临时生成一个更大的当前时间,假装它代表更新的业务事实。同版本却不同内容,也需要定义是拒绝、报警还是进行显式修正。
如果输出的是增量,而不是完整快照,则需要另一种身份:operation_id 应标识“一次应当执行的业务增量”。去重标记与增量生效必须在外部系统中原子协调。先登记“已经处理”再累加,中途失败可能漏算;先累加再登记,中途失败又可能重复算。
这里讨论的是根据写入语义推导出的设计约束,不意味着任意实现 upsert 的 Sink 都已经具备事务型 Exactly-once。连接器支持的写入模式、主键要求和故障保证,仍要逐项核对。3129
本章小结:Exactly-once 必须带着对象和边界来表达:对哪一段输入、哪一份状态、哪个外部结果,以及哪些故障成立。脱离这些限定,单独说“只执行一次”通常是不准确的。
九、故障恢复与作业演进:状态怎样跨越进程和版本
9.1 Restart Strategy 与 Failover Strategy
一个 Task 失败后,系统至少要回答两类不同问题。
Restart Strategy决定是否继续尝试运行,以及什么时候尝试。它可以表达重试间隔、退避方式和失败次数等约束。资源短暂不可用时,适当等待可能有意义;程序每次处理同一条记录都会抛出异常时,无限制快速重试只会重复同一个失败。
Failover Strategy决定哪些任务要被取消并参与恢复。只恢复发生异常的线程,不一定能重建与上下游一致的数据边界。Flink 的全局恢复与基于 Region 的恢复机制,处理的正是这个范围问题。32
恢复范围与执行图的连接方式有关。通过流水式数据交换连接的任务具有运行依赖,不能假设其中一个任务已经丢失状态,其他任务却可以毫无协调地继续使用旧执行过程中的数据。对于能够物化、保留并重新消费的结果分区,则可能在满足条件时复用它们,限制需要重算的范围。1332
这里必须避免一个误解:网络交换边界不等于故障恢复边界,keyBy 也不意味着故障能止步于这条边。 一个主要由流水式交换连接起来的流式作业,可能形成很大的恢复 Region;批处理中的阻塞式、可重用分区,才更容易成为阶段解耦的基础。
若依赖的中间结果不可用,恢复范围还可能进一步扩大。因此,“只挂了一个 TaskManager”不代表只需要恢复它上面的单个任务。
9.2 一次故障恢复的完整过程
以 TaskManager 丢失为例,一次恢复可以按如下顺序理解。
首先,协调端通过连接、心跳或执行失败报告识别异常,判断受影响执行范围,并根据重启策略决定下一次尝试。随后取消需要一起恢复的任务,为新的执行实例获得可用资源并重新部署。
接着,从可用的已完成检查点中选定恢复边界,将相应的算子状态分配给新的并行实例,恢复 Source 读取进度、定时器和其他参与协议的状态。使用非对齐检查点时,还要恢复必要的通道数据;输出侧则要恢复其写入或待提交状态。2432
之后,任务重新消费输入,重放检查点之后的数据,计算状态逐步追上故障发生前,再继续消化故障期间新产生的积压。这个过程依赖源端仍然保留所需数据,检查点文件仍然可用,并且程序、状态和连接器能够相互兼容。
大状态作业的恢复成本还取决于后端:有的需要读取较多快照文件来重建本地状态;有的采用远程状态和按需加载。可用的本地恢复数据可以加速,但不能把任务本地副本当作远程持久化恢复点的无条件替代。141817
必须区分两个时间点:任务重新进入 RUNNING,不等于业务延迟已经恢复正常。
假设恢复后积压为 B,实时输入速率为 λ,系统有效处理能力为 μ。在相同记录口径、速率近似稳定且 μ 大于 λ 的简化条件下,追赶耗时约为:
追赶时间 ≈ B / (μ − λ)例如积压 600 万条,每秒新增 10 万条,恢复后的处理能力为每秒 15 万条,净消化能力只有每秒 5 万条,追赶约需 120 秒;这还没有包含故障检测、调度和状态恢复时间。
如果 μ 小于或等于 λ,系统即使重启成功,也没有足够余量追上实时进度。因此,可恢复性不仅是“有没有快照”,还包括恢复后是否有能力消化积压。17

图 09:故障恢复与追赶——恢复运行和消化完积压是两个不同节点。
9.3 Checkpoint 与 Savepoint
Checkpoint 和 Savepoint 都可以保存作业状态,但它们服务于不同的运行目标。
Checkpoint 面向持续运行中的故障恢复,强调频繁创建和高效恢复。Savepoint 更面向计划性的升级、迁移、暂停和重新组织作业,强调可控生命周期与演进操作。这种区别不能简化成“一个自动、一个手动”:检查点也可以被显式触发和配置保留,保存点也可以由自动化运维流程创建。3334
| 维度 | Checkpoint | Savepoint |
|---|---|---|
| 主要目的 | 意外故障后的计算恢复 | 计划性变更、迁移与作业演进 |
| 常见格式 | 状态后端相关的原生格式,可采用增量形式 | 可使用规范格式,也存在原生格式选项 |
| 所有权与生命周期 | 通常由 Flink 管理,按配置保留或清理 | 由用户控制保留和删除 |
| 可迁移性 | 取决于快照类型、共享文件和兼容性支持 | 规范格式更注重迁移,但仍受状态及后端支持条件约束 |
| 是否保证任意新程序可读 | 否 | 同样不保证 |
规范格式与原生格式也各有含义。规范格式为状态后端切换等操作提供更好的通用表示;原生格式保留更多后端特征,可能更有利于创建与恢复速度,但不能把它理解成完全相同的跨后端能力。并非所有后端都支持全部保存点格式。1434
管理增量检查点时,要特别注意共享文件引用。恢复依赖可能不只是一眼看到的单个目录;不能因为“最新检查点已经生成”,就随意删除它仍然引用的状态文件。保存点通常围绕自包含、可迁移的运维需求设计,但移动和删除也仍应遵循对应工具与格式的规则。1633
还有一个计算之外的边界:从更老的 Savepoint 回退,不会自动把外部数据库、已提交事务或已经发送的消息回退到同一历史时刻。计划性停止、恢复与外部输出身份需要一起设计,不能只检查状态文件是否加载成功。
9.4 扩缩容与状态兼容性
并行度从 3 改为 4 时,Flink 需要重新分配第四章介绍的 Key Group。每个 Key 仍然应该找回属于自己的旧状态,只是持有该状态的并行实例可能变化。这里迁移的是计算归属,不是把旧进度清空后从头来过。110
跨程序版本恢复,则还要解决“这份状态现在对应谁、应该怎样读取”。
首先是算子身份。为需要长期演进的有状态算子设置稳定的 uid,可以减少作业图调整导致自动生成标识变化的风险。展示名称 .name(...) 与恢复身份 .uid(...) 不是同一个机制。15
其次是状态身份和序列化兼容性。状态描述符中的名称用于识别状态;快照中的字节需要由兼容的序列化逻辑读取。仅仅保持 Java 变量名或算子展示名称不变,并不足以维持恢复兼容性。935
例如,旧版本的累加器是:
// V1:结构示意,不代表完整类定义class Accumulator { public long count; public double sum;}新版本为了十进制精度,直接把 sum 改成 BigDecimal:
// V2:这种字段类型替换不能直接假设与旧状态兼容class Accumulator { public long count; public BigDecimal sum;}旧快照记录的是旧类型的序列化内容,不能因为 Java 源码能编译,就认为旧字节能按新类型直接解释。以文档中支持演进的 POJO 序列化规则为例,增加或移除字段与改变现有字段类型具有不同的兼容性结论。实际能否迁移取决于使用的序列化器和快照格式,而不是“都叫 Accumulator”。35
不兼容变更可能需要显式迁移状态、使用受支持的 State Processor API 路径转换快照,或从可重放输入重建新状态。具体方式还要满足对应后端和快照格式的支持范围。3435
最后,还要审视 Key 的语义。如果从 device_id 改成 (tenant_id, device_id),旧状态的分组身份已经改变;如果修改 Key 的序列化或最大并行度,也不能假设原状态会自然获得正确的新归属。保存了快照,只意味着保存了旧计算事实,不意味着新程序的任何解释都与它兼容。3515
本章小结:恢复不是简单地“重启进程并读取文件”,而是在新的执行位置上,使用兼容的程序重新接续同一份计算事实。
十、Flink SQL:一条持续查询背后的状态计算
10.1 动态表与持续查询
在 DataStream API 中,我们显式描述如何消费记录、维护状态和产生输出。Flink SQL 换了一个入口:把不断变化的数据看作动态表,用关系查询描述希望持续维护的结果。7
例如,设备测量事件不断追加到 measurements,执行下面的查询:
SELECT device_id, COUNT(*) AS measurement_countFROM measurementsGROUP BY device_id;这不是每收到一条数据,就扫描所有历史记录并重新执行一遍查询。运行时可以为每个设备维护一个计数器,收到新事件时增量更新对应状态。这种持续分组聚合本身就是有状态计算。36
对于 device-A,结果表可能依次包含:
某一时刻:device-A → 10下一时刻:device-A → 11第二个结果并不是“又来了一个数量为 11 的新设备”,而是同一结果行的计数从 10 变成了 11。
动态表是描述语义的模型,不意味着 Flink 必须把完整输入关系物化成一个供任意业务请求访问的数据库。运行时实际维护哪些状态,取决于查询计划:计数可能只需要累加器,Join 可能需要保存关联记录,排序和排名又需要其他结构。7
因此,一条持续查询的输出可能包括新增、更新和删除。输入只追加,也不代表输出一定只追加;分组计数就是最简单的反例。
10.2 Changelog 与输出系统的适配
将动态表交给下游时,需要把“表发生了什么变化”编码成记录流,这就是 Changelog 的职责。
追加流只表示新增行,适合结果本身不再改变的查询。撤回表达先移除旧结果的影响,再加入新结果。Upsert 表达则依赖唯一键,用新的行内容替换该键对应的旧内容,并用删除消息表达该键被删除。7
Flink 的行变化类型可以用 +I、-U、+U、-D 分别表示新增、更新前镜像、更新后镜像和删除。是否需要发送更新前镜像,取决于查询与下游所接受的 Changelog 模式,不是每个输出都必须同时携带四种消息。37
对于“device-A 的计数从 10 更新到 11”,关键在于下游是否把记录解释为同一结果行的变化:
| 表达方式 | 发送的内容 | 下游解释 | 最终结果 |
|---|---|---|---|
| 包含更新前后镜像的 Changelog | +I(A,10) → -U(A,10) → +U(A,11) | 撤回旧行,再加入新行 | A 的计数为 11 |
| 按稳定键 Upsert | UPSERT(A,10) → UPSERT(A,11) | 用同一键的新行替换旧行 | A 的计数为 11 |
| 错把完整结果当成独立计数 | 追加 10,再追加 11,然后求和 | 把两次结果快照累计起来 | 错误得到 21 |
前两种方式正确物化后,结果表仍然只有一行:device-A → 11。例如 Upsert Kafka 连接器正是围绕主键和更新、删除消息组织这种表变化表示。31
结果快照流不是增量流。 如果要发送增量,应明确表达成先贡献 10、再贡献 1;把完整结果 10 和 11 直接追加再求和,会改变原查询的含义。
同理,写入端声明了主键并不意味着输入一定满足主键约束,或者外部存储自动解决了全部乱序更新。必须让查询产生的变化模式、连接器要求和下游物化语义相互匹配。31
10.3 流式 Join 为什么需要状态
批量 Join 面对的是已知输入集合,而持续 Join 面对的是尚未到齐、还可能变化的数据。左侧先来的记录,要不要等待右侧未来的记录?已输出结果会不会因为右表更新而改变?这些语义决定了需要保存什么状态,以及什么时候能删除。
普通流式 Join:维护当前两张动态表之间的关联。
假设测量表与设备信息表按 device_id 关联,没有任何时间范围条件。某个设备信息很晚才到,仍然可能与此前的测量记录形成匹配;设备信息更新,也可能要求修改已有 Join 结果。为了增量维护这些关系,算子需要保留后续匹配和撤回所需的状态。对于不断追加且没有时间边界的输入,状态可能持续增长。38
不是所有普通 Join 都要保留所有历史版本:主键和更新语义会影响实际状态结构。但“没有时间边界,就天然可以删掉很久以前的数据”依然不成立。
区间 Join:只等待时间范围内的匹配。
例如,只关联同一设备前后五分钟内的测量与校准事件,其时间条件可以表达为:
m.event_time BETWEEN c.event_time - INTERVAL '5' MINUTE AND c.event_time + INTERVAL '5' MINUTE这段仅展示关联条件,完整查询还需要设备键等条件。与普通 Join 相比,时间范围限定了某条记录还可能等待哪些未来记录。对于事件时间区间 Join,运行时可以结合相应输入的 Watermark 推进,在已经不可能再产生合法匹配时清理状态。3823
这不等于“记录在机器内存里待满五分钟就删除”。事件时间与处理时间不同;另一个输入的进度没有推进时,未来匹配仍可能尚未到达。
事件时间 Temporal Join:选择事件发生时有效的版本。
假设设备校准系数会变化。一条发生在 12:00 的测量,应使用当时有效的系数,而不是处理它时数据库里最新的系数。这类需求可以用事件时间 Temporal Join,按事实记录的时间关联版本化维表。维表需要能够表达主键与版本时间,关联条件也必须满足主键等要求。38
这里不是等待“前后五分钟的所有系数”,而是寻找“12:00 时有效的那一版”。所需版本可能早于 12:00 很久,因此不能机械地把保留时间缩成一个很短的滑动区间。Watermark 推进后,可以清理不再需要的历史版本,同时保留仍可能用于后续匹配的有效信息。
查询外部系统当前值的 Lookup Join 又是另一种情况:外部值和缓存可能随着处理时刻变化,恢复重放时读到的结果未必相同。它不能自动替代按事件时间查询历史版本的语义。28
| Join 类型 | 主要等待或选择什么 | 清理状态的核心依据 |
|---|---|---|
| 普通流式 Join | 当前动态表中可能发生的匹配与更新 | 查询语义、键约束及明确的保留策略 |
| 事件时间区间 Join | 时间范围内可能到达的另一侧记录 | 关联区间与相关输入的时间进度 |
| 事件时间 Temporal Join | 事实发生时有效的维表版本 | 版本有效性、Watermark 与未来查询需要 |
SQL 状态保留配置,例如空闲状态 TTL,可以限制资源增长,但未必保持原本无界关系查询的完整结果。一条等待记录被提前清除后,未来再来匹配数据就可能无法得到本应存在的结果。保留策略应由业务允许的关联范围推导,而不是只根据“内存不够了”随意设置。39
10.4 将 SQL 映射回前面的运行机制
最后回到设备分钟平均值。
假设已经注册了只追加有效测量事件的 measurements 表,包含 device_id、measurement 和事件时间列 event_time。event_time 是声明了 Watermark 的时间属性,而不只是普通的时间戳字段。下面是建表语句中的时间字段定义片段,并非完整连接器 DDL:
-- 建表字段片段:完整 DDL 还需其他列和 WITH 连接器配置-- 示例时间字段采用 TIMESTAMP(3),输入解析应与其约定一致。event_time TIMESTAMP(3),WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND若源数据使用带时区的时间表达,应先明确如何归一化,或根据需求使用 TIMESTAMP_LTZ 与对应时区设置;不能把字符串截掉时区后直接认为业务时间没有改变。时间属性和 Watermark 的定义,是后续事件时间计算的前提。23
在上述表已正确注册的前提下,持续窗口查询可以写为:
SELECT device_id, window_start, window_end, SUM(measurement) AS measurement_sum, COUNT(*) AS measurement_count, AVG(measurement) AS measurement_avgFROM TABLE( TUMBLE( TABLE measurements, DESCRIPTOR(event_time), INTERVAL '1' MINUTE ))WHERE device_id IS NOT NULL AND measurement IS NOT NULLGROUP BY device_id, window_start, window_end;这里使用窗口表值函数 TUMBLE,为数据确定一分钟窗口,再按设备和窗口边界聚合。查询本身没有展开 State API,但窗口聚合的执行仍然需要状态和事件时间推进。40
沿着一条记录,可以把它映射回前面的机制。
进入 Source 与计划中的处理算子。 连接器读取事件,解析成表中的字段,投影和过滤等逻辑被优化进实际执行计划;SQL 中的一段表达式不必对应独立进程。
进入适当的分区。 属于同一聚合分组的数据必须汇集到能够共同维护该组结果的位置,执行计划因而可能包含按分组键进行的 Exchange。规划器也可能采用局部聚合加全局聚合,具体形态应查看物理计划,而不是从 SQL 字面逐词推断。36
更新窗口状态。 对一个设备的一个窗口,平均值可以通过 sum 和 count 增量维护。若存在多级聚合,合并的是具有正确代数意义的中间量;直接平均两个局部平均值,在两个分区记录数不同的时候一般是错误的。例如 1 条记录均值 10、9 条记录均值 20,总均值应为 19,而不是 15。
等待事件时间推进。 Watermark 表明窗口计算可以向前推进。本文这条标准的窗口聚合查询按照窗口最终输出语义工作,不像无窗口的持续分组聚合那样,每次新记录到来都必须把中间结果交给下游。40
这也需要与第五章的 DataStream 示例区分:那里显式设置了允许迟到,并讨论了首次输出之后的再次更新。这里不能因为也是“一分钟窗口”,就默认 SQL 自动继承相同的 allowedLateness 行为。对于超出窗口处理边界的迟到事件,应根据所选算子与业务补正机制另行设计。Watermark 延迟策略也不是迟到更新策略的同义词。1940
在 Checkpoint 中保存尚未结束的计算。 窗口没到输出时间,也可以被检查点保存;检查点完成,并不强迫窗口输出。故障时恢复相应的输入进度和窗口累计状态,之后继续处理尚未计入恢复边界的记录。24
将输出交给满足语义的 Sink。 对当前只输出窗口最终结果的查询,结果可以具有追加形态;但“逻辑上每个窗口一行”不等于失败重试时外部永远只写一次。第八章的事务提交或幂等身份约束仍然存在。4029

图 10:一分钟均值 SQL——分区、窗口状态、Watermark、Checkpoint 与 Sink 语义的运行时映射。
本章小结:SQL 隐藏了许多操作细节,但没有消除分布式执行、状态管理、事件时间和外部一致性问题。理解运行机制,才能判断查询的结果何时形成、为何更新,以及故障后怎样继续。
回到全文的核心问题。
数据持续到来时,Flink 用并行任务和数据分区组织计算,用托管状态保留历史贡献;事件乱序时,先通过事件时间确定业务归属,再用 Watermark 和窗口规则决定计算何时推进;下游变慢时,有限缓冲和背压限制过量生产,非对齐检查点则为拥堵条件下的快照提供另一种形成方式;执行节点失败时,一致的检查点让状态、输入进度和必要在途数据能够共同恢复。
但引擎不能替业务替换语义:状态保留多长时间、迟到事件是否补正、重复业务事件如何识别、外部结果如何提交、旧版本状态如何迁移,仍然必须被明确设计。
Flink 的核心不是让每条数据只经过一次代码,而是让持续变化的计算过程,在明确的时间、状态与故障边界内,始终能够被解释并正确接续。
参考资料
以下为正文脚注对应的 Apache Flink 官方文档与官方技术文章。版本相关说明固定使用 2.2 文档基线;连接器依赖与部署版本的兼容性需单独核对。网络栈文章仅用于说明通用机制,不将其历史默认配置视为当前版本配置。
Footnotes
-
Apache Flink, Stateful Stream Processing. ↩ ↩2 ↩3 ↩4 ↩5 ↩6 ↩7 ↩8 ↩9 ↩10
-
Apache Flink, Learn Flink: Fault Tolerance. ↩ ↩2 ↩3 ↩4
-
Apache Flink, Execution Mode (Batch/Streaming). ↩
-
Apache Flink, Data Sources. ↩ ↩2 ↩3
-
Apache Flink, DataStream Operators Overview. ↩ ↩2 ↩3 ↩4
-
Apache Flink, Dynamic Tables. ↩ ↩2 ↩3 ↩4
-
Apache Flink, Parallel Execution. ↩ ↩2 ↩3 ↩4
-
Apache Flink, A Deep-Dive into Flink’s Network Stack(2019-06-05). ↩ ↩2 ↩3
-
Apache Flink, Deployment Overview. ↩ ↩2
-
Apache Flink, Upgrading Applications and Flink Versions. ↩ ↩2 ↩3
-
Apache Flink, Tuning Checkpoints and Large State. ↩ ↩2 ↩3 ↩4 ↩5 ↩6 ↩7
-
Apache Flink, Disaggregated State Management. ↩ ↩2
-
Apache Flink, Process Function. ↩
-
Apache Flink, Timely Stream Processing. ↩ ↩2 ↩3
-
Apache Flink, Time Attributes. ↩ ↩2 ↩3
-
Apache Flink, Monitoring Checkpointing. ↩ ↩2
-
Apache Flink, Monitoring Back Pressure. ↩
-
Apache Flink, Checkpointing under Backpressure. ↩ ↩2 ↩3 ↩4 ↩5 ↩6 ↩7
-
Apache Flink, Determinism in Continuous Queries. ↩ ↩2
-
Apache Flink, Fault Tolerance Guarantees. ↩ ↩2 ↩3 ↩4
-
Apache Flink, Apache Kafka DataStream Connector. ↩ ↩2 ↩3
-
Apache Flink, Upsert Kafka. ↩ ↩2 ↩3
-
Apache Flink, Task Failure Recovery. ↩ ↩2 ↩3
-
Apache Flink, Savepoints. ↩ ↩2
-
Apache Flink, Checkpoints vs. Savepoints. ↩ ↩2 ↩3
-
Apache Flink, State Schema Evolution. ↩ ↩2 ↩3 ↩4
-
Apache Flink, Group Aggregation. ↩ ↩2
-
Apache Flink, DataStream API Integration. ↩
-
Apache Flink, Table Configuration. ↩
-
Apache Flink, Window Aggregation. ↩ ↩2 ↩3 ↩4