流处理架构:谁读哪个分区、状态放哪、哪个时钟说窗口完整、崩溃后从哪重放

比较 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 重建,或切到持有拷贝的 standbychangelog topic按 key 的计数、join、去重,且要快速接手没有 standby 时恢复随状态增长;换分区键要搬状态
Event-time Windows with Checkpoints按 key 分组的算子子任务存着算子状态和 source 位置的 checkpoint算子内部,由 checkpoint 快照事件时间、watermark、允许迟到、旁路输出恢复最近的 checkpoint 并回退 source;可回退 source 加事务 sink 才有 exactly-oncecheckpoint 存储有迟到数据、要求精确的时间窗口分析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。

面试时这样回答

  1. 先复述约束:要不要记住前一条、状态多大、结果按发生时间还是处理时间、几个团队消费、要回放多久以前。
  2. 说数据路径:点名图上的边,例如「分区分给 worker,写完 sink 再提交 offset」「更新本地状态并写 changelog,故障从 standby 接管」「watermark 越过窗口末尾就输出,迟到的进旁路,barrier 存快照后回退 source」「共享吞吐轮询、专属 fan-out 推送、新应用从序列号回放」。
  3. 说状态归谁:已提交的 offset、changelog topic、checkpoint 存储,还是留存的 shard 日志。
  4. 说代价与一个故障:例如空闲分区卡住 watermark 让所有窗口不输出,修法是空闲 source 检测把它排除、监控 watermark 停滞并告警。

一手证据