实时结果为什么可能需要修改
播放器在地铁里产生事件,网络恢复后才上传;服务端短暂重试让同一事件再次到达;负责统计的任务恰好在处理过程中重启。对业务而言,这些仍然是原来发生的播放行为。对计算系统而言,它们的到达顺序、处理时刻和可见结果已经不同。
如果直接按照服务器收到事件的时间分组,网络情况就可能改变事件属于哪个窗口。如果简单认为窗口第一次输出就是最终答案,后续迟到事件又无处可去。因此实时统计的核心,不只是“持续运行一段聚合代码”,还包括如何表达时间、维护状态和发布修订。
本篇以五分钟播放统计为例,使用 Flink 1.20.3、Kafka 输入与 Kafka 事务输出,实际观察重复、乱序、迟到和 TaskManager 重启。我们会明确接纳规则,并把实时结果与相同输入范围的参考答案比较。
三种时间回答不同问题
事件时间回答业务行为何时发生;到达时间回答记录何时进入某个接入边界;处理时间回答当前算子何时执行。它们可能相近,也可能相差很大。一次网络重试通常不会改变业务行为发生的时刻,却会改变记录被处理的时刻。
如果需求是统计某段业务时间内发生的播放,事件时间更接近指标语义;如果需求是观察服务器当前每秒处理多少条记录,处理时间本身就很有意义。不能只说事件时间更高级,而应先确定指标要描述哪个现实过程。
事件时间也需要可信来源。客户端时钟错误、时间单位混用、字段为空或未来时间异常,都可能干扰计算。生产系统应该定义合法范围和异常处理;本实验使用固定生成的 UTC 毫秒时间戳,将讨论集中到可控的乱序和迟到。
时间戳的表示与业务分组时区也应分开。原始事件保存 UTC,日报按上海时区转换;固定五分钟窗口按时间戳边界分组。跨自然日的统计要明确时区,不能让部署机器的系统时区无意决定结果。Flink 时间概念可作为进一步阅读。
为什么需要 Watermark
持续流没有一个自然的“全部到齐”信号。要输出事件时间窗口,系统需要对进度作出判断。Watermark 提供这种事件时间进度表达,使算子可以推进内部时间并触发相应动作。它不是外部世界已经完整上报的证明。
本实验使用有界乱序策略,依据已观察到的最大事件时间,保留五秒的乱序余量。这个五秒是教学配置,表达我们希望如何权衡及时输出与等待乱序,不是从真实网络统计推导出的服务承诺。
Watermark 推进之后,仍然可能收到更早的事件。这些记录是否进入结果,要看窗口和迟到策略。把“Watermark 已经过了某个时间”解释为“之后绝不可能出现这个时间之前的数据”,会让应用在真实延迟下错误丢数而不自知。
还要关注多输入。下游通常不能只看最快分区的进度,否则慢分区的数据会过早被判为迟到。一个长期没有事件的分区可能阻碍整体时间推进,空闲检测因此影响结果何时输出;而一个只是暂时变慢的分区与真正空闲的分区,又不完全相同。
本机实验对空闲输入使用三秒检测,并向各输入分区发送明确的教学时间推进记录。这些 tick 只推动样例的事件时间,不计入业务播放。真实系统不能随意伪造高时间戳事件来迫使窗口完成,应该根据源数据和业务完整性协议设计进度。
窗口定义了集合,触发与清理定义了生命周期
本系列使用五分钟、不重叠的事件时间窗口。某条事件归属哪个窗口,由它的事件时间决定,而不是由代码恰好在哪一秒执行决定。同一个窗口在首次触发后,还可能因接纳迟到记录而再次输出。
窗口长度、乱序余量和允许迟到是三个不同参数。窗口长度定义业务分组范围;乱序策略影响时间进度;允许迟到决定首次完成之后还保留状态多久,以便处理符合规则的晚到记录。不能把三者混成一个“延迟五分钟”。
实验允许窗口结束后三十秒范围内的迟到处理。超过清理边界的数据进入侧输出,保留事件身份供离线修正。侧输出不是把记录悄悄丢弃,也不意味着实时主结果已经包含了它。接纳与隔离应有可以检查的结果。Flink 窗口文档说明了窗口生命周期与迟到处理。
对外提供结果时,还应告诉使用者它是暂定值还是已按某个规则完成的值。实时窗口结束不等于所有源头数据永久到齐;最终业务确认可能要结合更长时间的离线补齐。这个区别会在第六篇的批流对账中直接体现。
状态是连续计算的工作记忆
统计一个窗口的播放次数与观看时长,需要保存已经累计到什么程度;识别重复事件,需要保存哪些 ID 已经出现。它们都属于计算状态。把状态保存在普通 Java 静态变量里,无法自动获得 Flink 的分区管理、快照与恢复能力。
示例先按事件 ID 分组,用托管状态判断是否已见;再按视频 ID 分组,在窗口中维护播放次数和观看时长。两次分组服务不同目的,前者管理身份,后者组织指标。输入 Kafka 的物理分区不能替代这两个业务状态键。
去重状态也需要有界。本例通过事件时间定时器在一天后清理已见 ID,避免无限保留。这个范围只是一项实验约定:超过去重保留范围的重复事件,不在本例同样的保证范围内。窗口状态和去重状态采用不同生命周期,也应分别解释。
状态越多,快照、恢复、扩缩容和内存管理的成本就越高。把允许迟到从三十秒改成三天,并不只是多等一会儿,可能让大量窗口长期保留。配置应该依据晚到分布、修正要求和资源预算,不能孤立地追求“尽量不丢任何迟到”。
Checkpoint 保存的是一致的恢复位置
任务失败时,需要同时知道输入读到了哪里以及状态累计到什么程度。如果只恢复计数,却从更早位置再次读入并累加,会重复计算;如果只恢复读取位置,却丢失此前状态,又会少算。恢复需要把这些部分放在一致的切面上理解。
Flink 的 Checkpoint 协调输入位置与托管状态的快照。故障后从已完成的 Checkpoint 恢复,重新执行尚未纳入该快照的部分。重新处理输入并不等于最终业务结果一定重复,是否产生重复外部效果还取决于输出端的协议。Checkpoint 文档介绍了基础机制。
本例启用两秒一次的 Checkpoint,将数据写入专用本地目录,配置有限次数的失败重启。实验中停止 TaskManager,再启动新进程,让任务通过恢复机制继续执行。这里没有重启 JobManager,也没有验证控制节点高可用或本地磁盘损坏后的恢复。
持久化目录本身的可靠性不可忽略。教学环境使用本机目录;多机器部署时,需要让替代执行节点能够访问恢复所需的状态。只看到 Checkpoint 成功,却把状态放在即将被销毁的临时磁盘上,无法满足真正的故障恢复目标。
输出一致性需要连接器参与
如果窗口每次输出都调用一个只支持累加的远程接口,即使内部状态恢复正确,也可能在重试时重复增加外部计数。因此需要继续追问:输出如何提交,消费者何时能看见,未完成事务如何处理?
本实验选择 KafkaSink 的精确一次交付模式,使用独立事务 ID 前缀,并由读取工具采用 read_committed。这样可以围绕 Kafka、Flink Checkpoint 和 Kafka 事务观察具体的端到端路径。Flink Kafka 连接器文档列出了使用条件。
事务超时必须容纳 Checkpoint 与恢复所需时间,不同作业的事务 ID 前缀也不能冲突。只设置一个枚举值而忽略这些条件,不能构成可靠输出。示例保留配置、作业身份和恢复记录,便于读者确认真正运行的是哪条路径。
精确一次也不等于业务唯一一次。上游发送两个不同 ID 表示同一播放行为,计算系统仍会将它们视为两个不同输入。这里能够验证的是:针对确定样例、明确身份和本次故障过程,状态与已提交结果满足我们的契约。
窗口更新必须按绝对值解释
示例输出键由窗口开始时间和视频 ID 组成,值包含当前窗口的播放次数与观看毫秒数。首次输出二,之后输出三,表示同一个窗口的最新绝对值由二改为三。消费者按同键采用较新的已提交记录,不能把两个值相加。
同键记录进入同一 Kafka 分区,读取工具按分区日志顺序处理更新,并生成最终窗口快照。若要把结果接到另一个数据库,可使用稳定键覆盖或带版本条件的更新,但具体一致性仍需要重新验证;本例没有实现任意外部数据库的事务写入。
榜单建立在最终窗口快照之上,按播放次数降序和视频 ID 升序排列。采用确定的并列规则,才能在不同引擎和不同运行中比较结果。如果两条记录并列而排序没有第二关键字,列表顺序变化未必代表计算错误。
业务还可能需要观察窗口每次修订,而不是只关心最新值。此时应保留更新历史,并明确版本、来源和发布时间。结果快照与修订事件流服务不同需求,不能因为都使用 JSON,就省略它们的语义说明。
最小实验:同一份样例经过真实链路
下载完整源码包,按运行说明启动隔离 Kafka 和 Flink,完成接入后执行:
python3 run.py streaming --work /absolute/path/to/lab-work
第一组实验读取一万条合法唯一事件、一百条重复和一条非法样例。推进事件时间并完成 Checkpoint 后,读取 Kafka 已提交输出。结果产生一百个视频窗口,与 Spark 离线聚合及独立参考答案逐组一致;诊断输出包含一百条重复记录和一条非法记录。
第二组采用四个业务事件和一个重复副本,按阶段发送。事件 a 和 b 先到,a 再次到达;等待 Checkpoint 完成后停止并重启 TaskManager。随后推进窗口、发送允许范围内的迟到事件 c,再推进到清理边界之外,最后发送过迟的事件 d。
| 阶段 | 本次实际观察 |
|---|---|
初始 a、b 和重复 a |
重复副本进入诊断输出 |
| TaskManager 停止并重启 | Checkpoint 恢复次数为 1 |
| 首次窗口输出 | 播放次数 2,观看时长 200 ms |
接纳迟到的 c |
同窗口更新为 3,观看时长 300 ms |
过迟的 d |
进入迟到侧输出,主窗口保持 3 |
如果离线计算包含 a、b、c、d 四条唯一合法事件,完整答案为四。实时主结果为三与这个完整答案不同,是实验预先定义的接纳规则造成的差异。第六篇会将侧输出和离线补数连接起来,而不会将三强行宣称为所有事件的最终真相。
排查“作业运行但没有结果”
首先确认输入确实到达作业。topic 名称、消费起点、分区发现和反序列化错误都会影响输入。仅看到 Job 状态为运行中,不足以证明窗口已经收到合法事件;应结合输入记录数、非法记录输出和消费位置检查。
然后观察事件时间进度。窗口尚未触发可能是最大时间戳还不够,也可能是一个分区没有推进或空闲判断没有生效。如果直接修改窗口大小来让结果出现,可能掩盖真正问题。应先明确当前 Watermark 和目标窗口边界。
再看 Checkpoint 和输出事务。结果已经进入未提交事务时,read_committed 消费者暂时看不到它是协议的一部分。持续无法提交,则需要检查 Checkpoint 是否失败、事务是否超时或任务是否反复恢复。可见性延迟与数据丢失不能混为一谈。
本次实施中也真实遇到过依赖版本不一致导致窗口序列化失败的问题,随后固定 Jackson 依赖并在打包时隔离其命名空间。这个错误发生在业务输出阶段,而不是输入阶段;只有检查失败日志和恢复次数,才能避免将“不断重启但仍显示任务存在”误判为运行正常。
升级和恢复也需要稳定身份
示例为输入、去重、窗口和事务输出设置稳定的算子 UID,帮助状态与逻辑节点建立明确对应。修改程序拓扑、状态结构或并行度时,不能仅因为新 JAR 编译成功就认为旧状态能够直接恢复。应先在隔离环境从目标快照启动,检查状态兼容和结果,再安排正式切换。Checkpoint 用于故障恢复的流程,与主动升级时保存、恢复状态的运维流程也应分别验证。
实验没有覆盖跨版本升级,因此固定所有组件与连接器版本。读者复现时先使用这套基线,再单独升级某个依赖,能更清楚地区分业务逻辑错误与兼容性变化。
实时系统的正确性是一份边界清晰的约定
本实验已经验证状态恢复、业务去重、窗口修订和过迟分流,但没有覆盖所有类型的故障。任务进程停止与 broker 故障不同,正常网络与网络分区不同,单机状态恢复与跨节点容灾也不同。每增加一个故障边界,都应补充对应输入、状态和输出的检查。
持续计算的价值,在于更早提供可以解释的结果,而不是隐藏数据还在变化的事实。将暂定值、迟到修订和离线校正分开表达,既方便用户理解,也为补数和审计留下空间。下一篇会把这些机制组合成一条能够发现差异、定位问题并修复历史的数据链路。