深入解析流处理系统两条核心约束——分区是排序单元,状态让重启变贵。覆盖 Kafka 分区与消费者并行度的关系,以及最常见的容量规划错误。
日志即架构
流处理系统是一个分区的、仅追加的日志,配有追踪其偏移量的消费者。用 Kafka 的术语来说,一个 topic 被分割为多个 partition;每个 partition 是一个有序的、不可变的序列;消息的 key 通过哈希决定其所属的 partition;消费者组将每个 partition 精确分配给一个消费者实例。这些设计后果直接由此而来,值得明确陈述,而非想当然。
顺序是按 partition 维度的,绝无全局顺序。两个具有不同 key 的事件,无论时间戳多么接近,都没有定义好的相对顺序。如果你的逻辑要求 user.created 在 user.updated 之前处理,两者必须使用相同的 key。
Partition 数量决定了消费者并行度的上限。消费者组中的实例数多于 partition 数时,会有空闲实例。这是容量规划中最常见的错误:有人将部署扩展到 32 个副本,却只有 12 个 partition,然后纳闷为什么吞吐量没有提升。
重分区并非免费的。增加 partition 数量会改变现有 key 的哈希映射,因此一个原本在 partition 3 上的 key 可能落到 partition 17 上,而该 key 的事件仍在 partition 3 上传输飞行中。任何你正在积累的带 key 的状态都会因此被分割到两个消费者。
热 key 是一个硬性天花板。一个 key 的流量无法分散。如果 40% 的事件携带相同的 tenant id,无论你创建多少个 partition,其中一个 partition 都承载着 40% 的负载。
这些保证的主要参考是 Apache Kafka 文档,Pulsar 和 Kinesis 在不同命名下有等价的概念——Kinesis 将 partition 称为 shard,并公布了 per-shard 写入上限,这使得热 key 问题变得显式。
一个完整的吞吐量计算
假设一个服务每秒产生 120,000 个事件,平均序列化大小为 800 字节,处理每个事件需要 1.2 毫秒的 CPU 时间。以下所有内容都来自这三个既定假设的算术运算;代入你自己的数字,答案的形状不会改变。
ingest bandwidth 120,000 ev/s x 800 B = 96 MB/s
x 3 replicas = 288 MB/s written to disk
CPU required 120,000 ev/s x 1.2 ms = 144 core-seconds per second
-> 144 cores fully busy
at 60% target utilisation -> 240 cores provisioned
partitions needed one consumer thread does 1/0.0012 = 833 ev/s
120,000 / 833 = 144 partitions minimum
with 2x headroom for skew -> ~288 partitions
per-partition rate 120,000 / 288 = 417 ev/s
417 x 800 B = 333 KB/s per partition
由此得出两个结论。Partition 数量由拓扑中最慢的阶段决定,而非由摄取速率决定:如果某个 enrichment 步骤因网络调用而需要 12 毫秒而非 1.2 毫秒,那么同样的流需要十倍的 partition,或者该步骤需要是异步的。60% 的目标利用率不是 padding——处于 95% CPU 的流处理器在任何暂停后都无法赶上,因此 lag 会单调增长,在不卸载负载或不增加容量的情况下永远不会恢复。
需要报警的数字是按时间衡量的消费者 lag,而非消息数。400,000 条消息的 lag 本身毫无意义;90 秒的 lag 精确告诉你每个下游答案的过期程度。计算方式为 lag-in-messages 除以当前消费速率,同时对导数而非绝对值报警——lag 在 30 秒处持平意味着系统处于平衡状态,而以每秒 2 秒速度增长的 lag 意味着到明天早晨系统将落后数小时。
状态才是难点
对流的无状态 map 映射既简单又罕见。任何进行聚合、连接、去重或模式检测的操作都持有按分区 key 索引的状态,而该状态必须在重启后存活。
标准设计将状态保存在每个任务本地的嵌入式键值存储中——Flink 和 Kafka Streams 都使用 RocksDB——并通过定期将一致快照写入对象存储,或将每次变更镜像到压缩的 changelog topic 来实现持久化。恢复时重新加载快照,并从快照记录的偏移量处重放日志。有两个数字决定了恢复的痛苦程度:checkpoint 间隔(决定了需要重放多少内容)和状态大小(决定了重新加载需要多长时间)。一个拥有 400 GB 状态和 5 分钟 checkpoint 间隔的作业,其恢复时间可能比导致它的故障本身还要长。
由此推论,无界状态是一种潜在的中断。一个连接保留它曾经见过的每个 key,或者一个没有过期时间的去重集合,可能完美运行六个月然后无法重启。每块带 key 的状态都需要生存时间(TTL),这在实践中意味着每个连接需要一个声明的窗口,每个去重需要一个声明的保留期——这与日志去重中相同的纪律如出一辙。
交付语义,精确版
At-most-once 在处理前提交偏移量:崩溃会丢失事件。At-least-once 在处理后提交:崩溃会重放事件。Exactly-once 并不是保证每个事件只处理一次——在不可靠的网络中这是不可能的——而是保证可观察效果只发生一次,通过将状态更新和偏移量提交变为单一原子事务来实现。
重要的限制是保证止步于事务系统的边界。写入 Kafka 的处理器可以是端到端 exactly-once,因为偏移量提交和输出写入join了一个事务。同一个处理器调用外部 HTTP API 则无法做到,因为该调用不在事务中,失败后会重试。唯一可靠的模式是使外部效果幂等,通常使用从事件衍生的确定性幂等 key,并接受 at-least-once 配合幂等 sink 作为实际的设计。
Lambda 架构对同一输入同时运行批处理路径和流处理路径,用流处理路径获得快速近似答案,再用批处理路径后续用精确答案覆盖。它诚实地承认了延迟数据的存在,代价是每条业务逻辑都以两种语言实现两次,并且会产生漂移。
Kappa 架构只保留流处理路径,通过重放来处理更正:保留源日志足够长时间,使得你可以从更早的偏移量重放到新的输出,然后切换。这更易于维护,并对保留期提出了硬性要求——你只能重放你保留的内容,因此源保留期现在是正确性参数,而非存储偏好。两者之间的权衡实际上是关于更正存在于哪里的权衡,在更简单的设置下同样问题在 batch 与 streaming 的对比中已有讨论。
无论你选择哪种形状,决定你的聚合实际含义的语义是窗口语义——滚动窗口、滑动窗口和会话窗口配合水印——这些值得在构建拓扑之前确定,而非之后。
Windowing Strategies for Streaming Aggregation
Change Data Capture for Feeding an AI Pipeline
Real-Time Fraud Detection From Event Streams