流处理架构:谁读哪个分区、状态放哪、哪个时钟说窗口完整、崩溃后从哪重放
比较 Stateless Consumer Group、Keyed State with Changelog、Event-time Windows with Checkpoints 与 Retained Shard Log with Independent Consumers 四种流处理拓扑:并行的单位是什么、处理进度记在哪、每个 key 的状态放在哪并怎么重建、按哪个时钟判断窗口完整以及迟到事件去哪、崩溃后从哪里重放、状态归谁,以及一次 worker 崩溃或一个卡住的 watermark 能波及多远、怎么收住。
事件源源不断地来,按分区排着,有时候乱序。处理器读它们,可能按 key 记状态,可能按时间开窗,再把结果写到某处。所有流处理方案都要回答四个问题:并行的单位是什么、一个分区读到哪了记在哪;每个 key 的状态放在哪、故障后怎么重建;哪个时钟决定一个窗口算完整、迟到的事件去哪;崩溃或新加入的消费者能从多久以前重放。这四个答案,而不是用的哪家产品,才定义了拓扑。
想亲手走一遍四种拓扑、注入故障再恢复,打开互动 Lab:/system-design-lab/stream-processing-architectures。
有约束的设计问题
一家公司先后要处理四条流:把事件过滤、补字段后路由到别的系统,每条独立处理但同一客户要有序;按用户累计计数和 join,状态很大,实例宕机后几秒内要接手,重复计数会影响计费;按事件发生的那一分钟统计指标,移动端离线缓存的事件会晚几分钟上传,结果要准确且不能重复写报表库;一条被五六个团队各自消费的流,新团队要能从几天前补数据,有的消费者对延迟敏感。每一条该怎么处理?
四种 topology signature
| 架构 | 并行单位 | 进度记在哪 | 每个 key 的状态在哪 | 时间与迟到 | 恢复 | State owner | 适合 | 主要代价 |
|---|---|---|---|---|---|---|---|---|
| Stateless Consumer Group | 一个分区归一个 worker | 写完 sink 后提交到 broker 的 offset | 没有:每条记录独立 | 只有处理时间 | 别的 worker 接管分区,从上次提交的 offset 重读 | 分区日志和已提交的 offset | 过滤、补字段、换格式、路由 | 崩溃后重投,sink 必须幂等;记录之间不能聚合 |
| Keyed State with Changelog | 每个输入分区一个任务 | offset 加每个状态存储的 changelog topic | 任务旁边的本地存储,镜像到 changelog | 事件时间加 grace period | 从 changelog 重建,或切到持有拷贝的 standby | changelog topic | 按 key 的计数、join、去重,且要快速接手 | 没有 standby 时恢复随状态增长;换分区键要搬状态 |
| Event-time Windows with Checkpoints | 按 key 分组的算子子任务 | 存着算子状态和 source 位置的 checkpoint | 算子内部,由 checkpoint 快照 | 事件时间、watermark、允许迟到、旁路输出 | 恢复最近的 checkpoint 并回退 source;可回退 source 加事务 sink 才有 exactly-once | checkpoint 存储 | 有迟到数据、要求精确的时间窗口分析 | watermark 用完整性换延迟;checkpoint 花时间和存储;sink 要事务或幂等 |
| Retained Shard Log with Independent Consumers | 每个应用各自读每个 shard | 每个应用自己的 per-shard checkpoint | 在各应用里,不在流里 | 各应用自定 | 任何应用从自己的 checkpoint 继续,留存期内可回放 | 带序列号的留存 shard 日志 | 很多团队独立消费一条流并重放历史 | 每 shard 读吞吐默认分摊,除非买专属 fan-out;留存期限制回放 |
批处理(有界输入、会结束的作业)是这些拓扑在需要持续结果时所替代的基线;exactly-once 是 changelog 和 checkpoint 拓扑只有搭配幂等或事务 sink 才成立的性质;流与流的 join 是两个输入上的 keyed state,属于第二种拓扑。它们在这里解释,不单独画。
1. Stateless Consumer Group:分区分给 worker,写完 sink 再提交 offset
Confluent 的 Kafka 介绍把地基写清了:「topic 被拆成分区,也就是一个 topic 的日志被拆成位于不同 broker 上的多个日志」;「事件 key 相同的事件,例如同一个客户 id 或车辆 id,被写进同一个分区,Kafka 保证该分区的任何消费者读到这些事件的顺序和写入顺序完全一致」;「事件在被消费之后不会被删除。topic 可以配置成在数据达到一定年龄或 topic 达到一定大小后过期」;「每个 topic 都可以复制,甚至跨地域或数据中心」。投递语义那一页给出三种保证的定义:「消费者读一批消息,先保存它在日志里的位置,再处理这批消息」是最多一次,「如果系统故障,消息可能丢失且不会被重新投递」;「消费者读一批消息,处理它们,再保存位置」是至少一次,「消息永远不会丢失,但可能被投递不止一次」;幂等生产者「保证重发一条消息不会在日志里产生重复条目,并且保持日志顺序」,因为「broker 给每个生产者分配一个 ID,并用生产者随每条消息发送的序列号去重」。
典型事故:为了省一次往返,代码先提交 offset 再批量写 sink,部署时 Worker A 在两步之间被杀——接手的 worker 从提交处往后读,那一批记录永远不会被处理,没有任何报错;改成先写后提交、接受重投、sink 按记录 id 去重,已经丢的那批按 offset 范围从 topic 重放。key 是商户 id,一个头部商户占一半流量,它所在的分区积压几百万条——加 worker 没用,一个分区只能一个 worker 读;给热 key 加后缀摊开,只在同一子 key 内保证顺序。
2. Keyed State with Changelog:状态放在任务旁边,每次更新也写进 changelog
Kafka Streams 的架构文档:「每个 stream partition 是一个全序的记录序列,映射到一个 Kafka topic 分区」;「stream partition 到 stream task 的分配从不改变」;「Kafka Streams 应用里的每个 stream task 可以嵌入一个或多个本地状态存储,通过 API 存取处理所需的数据」;「对每个状态存储,它维护一个有副本的 changelog Kafka topic,记录所有状态更新」;故障时「Kafka Streams 保证在恢复处理之前,通过回放对应的 changelog topic 把关联的状态存储恢复到故障前的内容」;「为了最小化恢复时间,可以给应用配置本地状态的 standby 副本,也就是状态的完整复制」,运行时「把任务分配给已经存在 standby 副本的实例」。概念页补上保证和时间:默认至少一次下「记录永远不会丢失但可能被重新投递」,设置 processing.guarantee='exactly_once_v2' 启用 exactly-once;事件时间是「事件或记录发生(即被来源创建)的时间点」,处理时间是「事件或记录恰好被流处理应用处理的时间点」,摄入时间是「事件或记录被 Kafka broker 存进 topic 分区的时间点」;grace period「控制 Kafka Streams 为一个窗口等待乱序记录多久」并「指示窗口结果何时最终」。read_committed 下消费者「只读来自已提交事务的消息」。
典型事故:状态几千万个 key 没有 standby,实例宕机后新实例必须从 changelog 头回放,几十分钟里分区 0 的事件全部积压——配置 standby 副本,故障时把任务派到有拷贝的实例只补最后一小段。默认至少一次,输出了新计数后、提交 offset 前实例被杀,同一事件再算一次——开 exactly-once,让状态更新、输出记录和消费的 offset 在同一个事务里提交。
3. Event-time Windows with Checkpoints:按事件时间开窗,watermark 说完整,checkpoint 存快照
Flink 关于时间:「处理时间指执行相应操作的机器的系统时间」;「事件时间是每个事件在其产生设备上发生的时间」;watermark 声明「该流的事件时间已经到达时间 t,意味着该流不应再有时间戳 t' <= t 的元素」;但「某些元素可能被任意延迟,使得不可能指定一个时间点,到那时某个事件时间戳的所有元素都已经出现」。窗口页:「滚动窗口大小固定且不重叠」;「如果滑动步长小于窗口大小,滑动窗口可以重叠」;会话窗口「在一段时间没有收到元素时关闭」;「在 watermark 越过窗口末尾之后、但在越过窗口末尾加允许迟到之前到达的元素,仍然会被加进窗口」,触发「迟到触发」;Flink 保留窗口「直到允许迟到过期」,然后「移除窗口并删除其状态」;应用可以通过旁路输出「获得被当作迟到而丢弃的数据流」。Beam 用同样的三种窗口,再加上事件时间、处理时间和数据驱动三类触发器,以及累积与丢弃两种累积模式。
状态和恢复:「keyed state 可以看作一个嵌入式的键值存储」;「barrier 永远不会超过记录,它们严格顺着流走」;「一旦最后一个输入流收到了 barrier n,算子发出所有待发的记录,然后自己发出快照 n 的 barrier」;故障时「系统重启算子并把它们重置到最近一次成功的 checkpoint」,重放的记录「保证没有影响之前已经 checkpoint 的状态」;跳过对齐得到的是至少一次,因为「算子会继续处理所有输入,即使 checkpoint n 的某些 barrier 已经到达」;端到端 exactly-once 需要可回退的 source,「Apache Kafka 有这个能力,Flink 的 Kafka 连接器利用了它」。
典型事故:source 的八个分区里有一个来自已下线地区,几小时没有事件,算子取最小 watermark 停在几小时前,所有窗口都不输出——启用空闲 source 检测,超时的分区暂不参与 watermark 计算。移动端事件晚几分钟上传,而允许迟到设为零、没有旁路输出——每分钟的计数系统性偏低且不可追溯;按观测到的迟到分布设允许迟到,超过的送旁路输出并计数。
4. Retained Shard Log with Independent Consumers:一条留存的日志,多个应用各自读、各自回放
Kinesis 的术语页:「一个 Kinesis 数据流是一组 shard。每个 shard 有一个数据记录序列。每条数据记录有一个由 Kinesis Data Streams 分配的序列号」;partition key「用来在流内按 shard 对数据分组」,通过「MD5 哈希函数」映射;「每条记录的序列号在其 shard 内按 partition key 唯一」且「同一 partition key 的序列号通常随时间递增」;留存期「创建后默认设为 24 小时」,「可以把留存期增加到最多 8760 小时(365 天)」;「一条流可以有多个应用,每个应用可以独立、并发地消费流中的数据」;客户端库「确保每个 shard 都有一个记录处理器在运行」并「使用 Amazon DynamoDB 表存储与数据消费相关的元数据」。fan-out 页把两种消费者放在一起:共享吞吐消费者的读吞吐「固定为每 shard 总共 2 MB/秒。如果多个消费者读同一个 shard,它们共享这个吞吐」;专属 fan-out 下「每个注册使用 enhanced fan-out 的消费者获得自己的每 shard 读吞吐,最高 2 MB/秒,独立于其他消费者」,并且「Kinesis Data Streams 通过 HTTP/2 用 SubscribeToShard 把记录推送给你」,而共享消费者是「通过 HTTP 用 GetRecords 的拉取模型」;每条流可注册的专属消费者数量有上限。
典型事故:App A 的部署坏了两天没人发现,留存期是默认的一天——它恢复后 checkpoint 之后的一整天记录已经过期,一整天的事件对它永久丢失;留存期要覆盖最长故障,监控每个应用的滞后并在接近留存期时告警,把流归档到对象存储作为留存期之外的回放来源。五个团队各自用轮询消费者读同一条流——每个消费者的延迟成倍上升,最重要的和最不重要的一样慢;给需要独立吞吐的消费者注册专属 fan-out。
四种拓扑共同的底线
- 说清状态归谁。 带 offset 的分区日志、changelog topic、checkpoint 存储,还是留存的 shard 日志。
- 说清端到端的保证。 至少一次要配幂等 sink;exactly-once 只在可回退的 source 加事务或幂等 sink 时才成立。
- 说清哪个时钟决定完整。 处理时间只对无状态转换够用;按发生时间统计就要 watermark、允许迟到和旁路输出。
- 先记进度再动手。 offset 在写完 sink 之后提交,checkpoint 连同 source 位置一起存,changelog 在本地更新的同时写。
- 让滞后和留存对得上。 分区滞后、standby 滞后、watermark 停滞、消费者滞后对留存期,都要有指标和告警。
故障与恢复
| 架构 | 故障 | 用户看到什么 | 恢复 |
|---|---|---|---|
| Consumer Group | 先提交 offset 再写 sink,中途崩溃 | 一批记录消失 | 先写后提交;sink 按 id 去重;按 offset 范围重放 |
| Consumer Group | 写 sink 后、提交前崩溃 | 记录写两次 | sink 按记录 id 去重 |
| Consumer Group | 热 key 压满一个分区 | 一个 worker 累死其他闲着 | 给 key 加后缀或换高基数 key;按分区监控滞后 |
| Keyed Changelog | 大状态没有 standby,任务挂了 | 分区停滞到回放结束 | standby 副本;任务派到有拷贝的实例 |
| Keyed Changelog | 至少一次下崩溃重算 | 计数多一 | exactly-once:状态、输出、offset 同一事务;read_committed |
| Event-time Windows | 空闲分区卡住 watermark | 所有窗口不输出 | 空闲 source 检测;watermark 停滞告警 |
| Event-time Windows | 允许迟到为零 | 迟到事件无声消失 | 按迟到分布设允许迟到;旁路输出并计数 |
| Event-time Windows | 恢复后重放进非事务 sink | 结果重复 | 事务 sink,或以窗口和 key 为键幂等写 |
| Shard Log | 消费者停得比留存期久 | 一段记录永久错过 | 留存覆盖故障;滞后对留存告警;专属 fan-out;归档 |
怎么选
- 逐条转换不需要记忆:Stateless Consumer Group,写完 sink 再提交 offset,sink 幂等,按 key 分区只依赖同 key 内的顺序,分区数按最慢的 worker 定。
- 按 key 聚合、join、去重且要快速接手:Keyed State with Changelog,每个状态存储配有副本的 changelog,大状态配 standby,分区键一次定好,重复计数不可接受就开 exactly-once。
- 结果取决于事件何时发生且迟到真实存在:Event-time Windows with Checkpoints,在 source 处贴时间戳和 watermark,允许迟到有意设置,旁路输出有人看,checkpoint 间隔大于 checkpoint 时长,可回退 source 配事务或幂等 sink。
- 很多团队读同一条流并要重放历史:Retained Shard Log,留存期覆盖最长故障和回放需求,不能分摊吞吐的消费者注册专属 fan-out,每个应用的 checkpoint 存在自己的持久表里。
- 无论哪种:写下谁拥有进度、哪个时钟决定完整、端到端的保证到底是什么。
回到开头的四条流:过滤路由走 Consumer Group;按用户累计走 Keyed State with Changelog;每分钟指标走 Event-time Windows;多团队共享的流走 Retained Shard Log。
面试时这样回答
- 先复述约束:要不要记住前一条、状态多大、结果按发生时间还是处理时间、几个团队消费、要回放多久以前。
- 说数据路径:点名图上的边,例如「分区分给 worker,写完 sink 再提交 offset」「更新本地状态并写 changelog,故障从 standby 接管」「watermark 越过窗口末尾就输出,迟到的进旁路,barrier 存快照后回退 source」「共享吞吐轮询、专属 fan-out 推送、新应用从序列号回放」。
- 说状态归谁:已提交的 offset、changelog topic、checkpoint 存储,还是留存的 shard 日志。
- 说代价与一个故障:例如空闲分区卡住 watermark 让所有窗口不输出,修法是空闲 source 检测把它排除、监控 watermark 停滞并告警。
一手证据
- Confluent:Apache Kafka introduction
- Confluent:Message Delivery Guarantees(Kafka design)
- Confluent:Kafka Streams Concepts
- Confluent:Kafka Streams Architecture
- Apache Flink:Timely Stream Processing
- Apache Flink:Windows
- Apache Flink:Stateful Stream Processing
- Apache Beam:Programming Guide — Windowing
- Amazon Kinesis Data Streams:Terminology and concepts
- Amazon Kinesis Data Streams:Enhanced fan-out consumers