:流处理)
本文是《Designing Data-Intensive Applications》DDIA中文译名《数据密集型应用系统设计》第 11 章的导读。DDIA 是 Martin Kleppmann 所著的分布式系统经典本系列逐章导读把书的核心概念讲清楚。一句话主旨流处理把攒批再算变成来一条算一条——数据不再是静止的有界数据集而是持续到来的无界事件流。这带来一个根本难题事件到达的顺序和时间不可控所以流处理必须自己定义什么时候算到齐了——这就是水位线watermark和窗口window。本章的核心是时间语义事件时间 vs 处理时间的分离是流处理区别于批处理的灵魂。核心概念拆解1. 消息代理Message Broker——流的载体生产者(产生事件)Topic/Queue(消息代理: Kafka/RabbitMQ)消费者1(流处理)消费者2两种消息模型队列Queue发布订阅Pub/Sub代表RabbitMQKafka分发每条消息给一个消费者负载均衡每条消息给所有订阅者持久消费后即删按保留期存可重放适用任务分发消息触发动作多下游各取所需关键区别——消息可重放性RabbitMQ消费即删消息不可重放。消费者崩了没消费的消息还在但已消费的没了。Kafka消息按日志持久化消费者维护 offset。可以回退 offset 重放历史消息——这是流处理容错的基础批处理能重跑是因为输入有界且不变流处理能重放是因为 Kafka 把消息当日志存。消息代理选型看要不要重放消息触发动作不需要重放用队列模型RabbitMQ要重放历史或消费失败要回放用发布订阅Kafka。2. 事件时间 vs 处理时间——流处理的灵魂这是全章最重要的概念。网络延迟/乱序, 两者不一致事件时间 Event Time事件实际发生的时间(数据自带的timestamp)处理时间 Processing Time流处理系统收到事件的时间(系统时钟)为什么分离事件在源头产生事件时间 t1经过网络传输、队列排队、系统调度到达处理系统处理时间 t2。t2 t1且差值不定。网络乱序事件时间早的可能晚到t110:00 的事件比 t110:01 的晚到处理系统。如果按处理时间算10:00~10:01 的数据10:00 的事件可能还没到算不全。流处理的正确做法按事件时间分窗“我要 10:00~10:01 真实发生的事件”而不是按处理时间。但这要求等迟到的事件到齐——等多久这就是水位线要回答的。水位线 我认为事件时间 ≤ T 的事件都到齐了的承诺。水位线前进到 T就触发 T 之前的窗口计算。水位线本质是在延迟和正确性之间做权衡——水位线激进等得短结果出得快但可能漏迟到事件水位线保守等得久结果准但延迟高。3. 水位线Watermark——“什么时候算到齐”窗口已触发, 走迟到处理(丢弃/侧输出/重触发)事件时间 t10:00 到达, 水位线09:58事件时间 t09:59 到达(乱序,非迟到), 水位线09:58事件时间 t10:02 到达, 水位线推进到10:00水位线10:00 → 触发09:59~10:00窗口计算事件时间 t09:55 到达(迟到: 水位线已过09:55)水位线的生成策略固定延迟水位线 最大事件时间 − Δ允许 Δ 秒延迟。简单常用。百分位按事件延迟分布动态调整更智能但复杂。水位线过后的迟到事件丢弃最简单容忍少量丢侧输出送到单独流另行处理允许延迟重新触发窗口最准但最复杂水位线只对流式、按事件时间分窗的场景有意义。批处理攒文件一次性处理没有事件流的概念没有到齐的要求——文件是完整的。不是所有数据处理都要水位线。问题→方案问题——数据持续到来要低延迟处理而非攒批。场景——无界事件流事件到达顺序和时间不可控网络延迟致乱序、早事件晚到。方案——分离事件时间与处理时间按事件时间分窗计算用水位线承诺事件时间 ≤ T 的事件都到齐了来触发窗口水位线 延迟与正确性的权衡等得短出得快但漏迟到事件等得久结果准但延迟高。Kafka 把消息当持久日志存可回放 offset作为容错基础。流表对偶表的变更日志是流流回放物化成表。4. 窗口Window——把无界流切成有界块无界流无法算总计必须切成有界窗口分别计算。三种窗口固定窗口 Tumbling10:00~10:05, 10:05~10:10不重叠跳跃窗口 Hopping10:00~10:05, 10:02~10:07重叠(按步长前进)会话窗口 Session按活跃度切, 有间隔才断不固定固定窗口等长不重叠。最简单适合周期统计。跳跃窗口Hopping等长、按固定步长前进可重叠每 2 分钟算一次过去 5 分钟。适合移动平均。注部分框架把这种等长可重叠窗口叫滑动窗口但 Flink 等另一些框架中滑动窗口指按每条事件滑动的窗口术语存在冲突。会话窗口按活动间隙切用户操作间隔超时则断开窗口。适合用户行为分析。5. 流表对偶——流和表是同一事物的两面DDIA 第 10 章一个深刻观点流和表是对偶的。表的变更日志回放日志得到表表 Table当前状态流 Stream状态变化的日志表 → 流表的每次变更insert/update/delete是一个事件变更序列就是流CDCChange Data Capture。流 → 表把流的事件从头回放物化成当前状态就是表。实践意义数据库的 binlog 就是流物化视图就是流物化成表流处理持续更新一个物化视图视图是表更新源是流去重不是删除旧行是流回放时只保留最新版本——日志压缩按 key 只留最新值去重是流到表的聚合逻辑6. 流处理 vs 批处理 vs 微批——延迟谱系批处理攒足再算, 延迟分钟~小时吞吐最高微批 Micro-batch切小批, 延迟秒~分(Spark Streaming)流处理逐条算, 延迟毫秒~秒(Flink)批微批流延迟分钟~小时秒~分毫秒~秒吞吐最高中看引擎复杂度低中高时间语义/状态/水位线容错重跑作业重跑微批检查点状态恢复流处理的复杂度水位线/状态/检查点/乱序远高于批只在延迟要求 批能给的时才值得。如果场景不需要实时批处理的吞吐优势远超流的低延迟上流是过度工程。Mermaid流处理时间语义时间语义(核心)事件时间数据自带, 不可控顺序水位线事件时间到齐的承诺窗口按事件时间切有界块处理时间系统时钟, 随到达而定下一篇第 12 章——数据系统的未来。把批和流统一起来看派生数据如何把多个系统连成一致的全链路。