ARTICLE / 2 MIN READ
Kafka 加 Flink 构建实时数仓的关键细节
讲清事件时间、窗口、状态、乱序和 Exactly-once,避免实时指标看似及时却无法对账。
实时计算的难点不是把消息读出来,而是在延迟、乱序、重启和重复投递同时存在时,仍然得到可解释的结果。
事件时间优先
处理时间是消息到达计算节点的时间,事件时间是业务真正发生的时间。跨网络、客户端离线和重试都会造成到达顺序与发生顺序不同,统计订单和支付时应优先使用事件时间。
水位线表示系统认为某个时间点以前的事件大部分已经到达。允许迟到时间要结合业务延迟预算设置,不能为了“绝对准确”无限等待。
窗口与状态
滚动窗口适合固定周期指标,滑动窗口适合趋势观察,会增加计算量,会话窗口适合按用户行为间隔聚合。窗口状态要设置 TTL,避免异常用户或高基数 Key 让状态无限增长。
事件 -> 按用户/商品 KeyBy
-> 事件时间 + 水位线
-> 窗口聚合
-> 结果表/消息
Exactly-once 的边界
Flink 的状态一致性依赖 checkpoint,端到端 Exactly-once 还要求 Source 可重放、Sink 支持事务或幂等。把结果直接写入不支持幂等的外部 HTTP API,不能宣称端到端 Exactly-once。
常见做法是写入支持事务的消息或数据库,再由下游按事件 ID 去重。checkpoint 间隔、超时和保留数量要根据状态大小和恢复时间预算调整。
乱序和更正
窗口关闭后到达的迟到事件,可以丢弃、写入侧输出或触发结果更正。对财务指标不能静默丢弃,应该记录迟到量并进入补算流程。
运行监控
重点看消费位点延迟、每个算子反压、checkpoint 成功率、状态大小、窗口迟到数、输出重复率和恢复耗时。任务绿灯但消费延迟持续上涨,说明实时链路已经失去时效。
实时数仓要把“快”和“准”的取舍写进数据契约,并为每个指标保留重算入口,否则出了乱序或重启问题就无法解释结果。
GITHUB DISCUSSION
评论与回复
评论保存在 GitHub Discussions,登录 GitHub 后即可参与,发布和回复都在本页完成。
评论区进入视口后自动加载。