编程 Kafka Diskless Topics 深度拆解:当消息队列决定把「盘」从架构里删掉——WAL Segment、Diskless Coordinator 与 SQLite 状态机的一次成本手术

2026-08-09 01:48:03 +0800 CST views 8

Kafka Diskless Topics 深度拆解:当消息队列决定把「盘」从架构里删掉——WAL Segment、Diskless Coordinator 与 SQLite 状态机的一次成本手术

一、开场:你的 Kafka 账单里,最贵的那一项可能不是机器

先说一个很多团队都算过、但很少算明白的账。

假设你在 AWS 上跑一个三可用区的 Kafka 集群,replication.factor=3min.insync.replicas=2,生产者写入速率 100 MiB/s。你以为花钱的地方是 EC2 实例和 EBS 卷,实际上真正吃掉预算的是这三块:

  1. 副本复制的跨 AZ 流量。一条消息写到 leader,leader 在 AZ-a,两个 follower 分别在 AZ-b、AZ-c。于是每写入 1 字节,就要跨 AZ 传输 2 字节。AWS 对同 Region 跨 AZ 流量双向计费,官方价格是 $0.02/GiB(进出各算一次的口径要按实际账单确认,但量级如此)。
  2. 生产者/消费者的跨 AZ 流量。生产者在 AZ-a,分区 leader 在 AZ-c,写请求必须打到 leader——又是一次跨 AZ。消费者同理,Kafka 虽然有 KIP-392 的 follower fetch(rack-aware 读),但写路径没有对应机制,因为写必须去 leader。
  3. 块存储。EBS gp3 要为三份副本各存一份,容量 ×3,IOPS 和吞吐还要单独买。对比 S3 的单价,差着一个数量级。

把这三项摊开算一算:100 MiB/s 的入口流量,一个月大约是 100 × 3600 × 24 × 30 / 1024 = 约 253 TiB。复制放大 2 倍 → 506 TiB 跨 AZ;假设生产者有 2/3 概率不在 leader 所在 AZ → 约 169 TiB;消费者按 2 个消费组、rack-aware 部分生效算,再加几十 TiB。合计跨 AZ 流量轻松突破 700 TiB/月,乘以 $0.02/GiB,光网络费就是 一万四千美金以上,而且这部分钱不产生任何计算价值——纯粹是为了「把同一份数据在机房内部搬三遍」。

这就是 2026 年 Apache Kafka 社区推动 KIP-1150: Diskless Topics 的直接动因。这个 KIP 现在的状态是 Accepted(2026 年 3 月初仍在持续更新),它的两个落地子提案 KIP-1163: Diskless CoreKIP-1164: Diskless Coordinator 处于 Under Discussion,作者阵容基本是 Aiven 那批人加上社区核心 committer。

KIP-1150 的原文里有一句话,我觉得是整个提案最诚实的表述:

Apache Kafka is designed around low-durability block storage and direct replication.

Kafka 是围绕「低持久性块设备 + 自己做复制」设计的。这个前提在 2011 年的 LinkedIn 机房里成立——那时候你手上只有一堆机械盘,想要不丢数据只能自己写三份。但在 2026 年的云上,你手上有一个持久性 11 个 9、天然跨 AZ、还自带纠删码的对象存储。继续让 Kafka 自己做复制,本质上是为一个已经不存在的约束付费

这篇文章会把 Diskless 这套东西从头拆到尾:数据面怎么改、元数据面怎么改、offset 的全局序谁来定、幂等和事务怎么活下来、延迟涨到多少、账到底能省多少、以及在你决定上车之前应该知道的那些坑。


二、为什么 Tiered Storage(KIP-405)不够

先解决一个高频误解:「我们已经开了分层存储把冷数据放 S3 了,是不是就已经 diskless 了?」

不是。而且差得很远。

KIP-405 的分层存储做的事情是:把已经封存的 inactive segment 搬到对象存储。它的数据流是这样的:

Producer → Leader 本地盘(active segment)
             ↓ 复制(跨 AZ!)
          Follower 本地盘 ×2
             ↓ segment roll 之后
          上传到 S3,本地删除

注意关键点:active segment 依然是三副本写本地盘,依然要跨 AZ 复制。分层存储省的是「长期保留的容量成本」,一分钱的复制流量都没省。

而现实中,绝大多数 Kafka 集群的保留期是 3 天到 7 天。在这个保留期下:

  • 容量成本本来就不算大头;
  • 跨 AZ 流量是按写入量线性增长的,跟保留期完全无关。

所以分层存储解决的是次要矛盾。KIP-1150 里说得很直白:Tiered Storage "does not remove the need for replication and durable storage of active segments, which is the most substantial infrastructure cost"。

Diskless 要动的是主要矛盾:干掉直接复制,让写路径的落盘持久化交给对象存储。

顺带澄清一个命名问题——KIP 原文专门写了一节解释「Diskless 到底还用不用盘」:

Diskless is to "No Disks" as Serverless is to "No Servers."

盘还在,只是盘不再是 user data 的 source of truth。本地盘依然会用来:

  • 存 KRaft 元数据(这个跑不掉);
  • 存 batch metadata(取决于 coordinator 实现,KIP-1164 用的就是本地 SQLite + 内部 topic);
  • 在上传到 tiered storage 之前临时落地;
  • 做 consumer 读取的缓存。

丢了这些本地数据,不会丢用户数据,重建就是了。这是 diskless 和 classic 最本质的语义差别。


三、名词表:先把概念对齐

Diskless 引入了一堆新词,混淆了后面全看不懂。按 KIP-1163 的 Glossary 整理一下:

术语含义类比
Classic Topic现状:直接 append 本地块设备 + broker 间直接复制你现在跑的所有 topic
Tiered TopicClassic + KIP-405,inactive segment 下沉对象存储加了个冷存档
Diskless Topic不直接 append 块设备、不做直接复制本文主角
WAL Segment对象存储上的一个对象,混装多个 topic-partition 的 batch一个「合租」的 append-only 文件
Batch Coordinate「某个 batch 在某个 WAL Segment 的某个字节区间」的引用指针:(objectKey, offset, length)
Diskless Coordinator (DC)batch coordinate 和 WAL Segment 的唯一真相源,负责定全局序、分配 offset「排号机」
Diskless-Ready Cluster同时能承载 classic / tiered / diskless 三种 topic 的集群混合部署
Rack拓扑提示,云上通常等价于 AZ计费边界

一句话概括整个架构:

数据走对象存储(无序、可并发、便宜),元数据走 Kafka 内部 topic(有序、强一致、贵但量小)。

这个「数据面和元数据面分离」的思路不新——WarpStream、AutoMQ、Confluent Freight 都是这么干的。KIP-1150 的 Motivation 里也大方承认了:

Multiple protocol-compatible alternatives to Apache Kafka now use object storage... These alternatives are finding market success and their adoption is rising.

翻译过来就是:再不做,市场就没我们什么事了。KIP 里甚至直接写了一条收益是 "avoid obsolescence"(避免被淘汰)。一个 Apache 项目的设计文档能写这句话,说明社区是真着急了。


四、架构拆解(一):Produce 路径为什么能「任意 broker 写」

4.1 Classic 的写路径为什么必须去 leader

先复习一下 classic topic 的 Produce 处理,broker 要干六件事:

  1. 校验数据(CRC、magic、压缩格式、record 数量);
  2. 分配 offset 和 timestamp
  3. 把 offset/timestamp 注入到 batch 二进制里
  4. 写入持久化存储;
  5. 等待副本复制完成(acks=all);
  6. 返回响应。

第 2 步是命门。offset 是一个分区内严格递增、无空洞的整数,它必须由一个单点分配,否则两个 broker 同时写同一个分区就会撞号。Kafka 的做法是:这个单点就是 leader,leader 自己维护 nextOffset,单线程 append。

所以「写必须去 leader」不是网络层的限制,是 offset 全局序 的限制。

4.2 Diskless 怎么破这个局

KIP-1163 的关键动作是:把「分配 offset」和「写数据」拆开

  • 写数据:任意 broker 都能干。因为写到对象存储的数据此刻还没有 offset,它只是一坨字节,没有顺序含义。
  • 分配 offset:交给 Diskless Coordinator。broker 把「我写了哪些 batch、在哪个对象的哪个字节区间」这份元数据提交给 DC,DC 单点定序。

原文里这段描述很精确:

The data being persisted to object storage has no implicit order, and ordering is delegated to the batch coordinates. Coordinates passed from each broker to the diskless coordinator are locally ordered... The diskless coordinator performs the append and assigns a global order.

也就是说:同一个 broker 提交的 coordinate 之间有序,跨 broker 无序,由 DC 归并成全局序。

于是六件事的职责重新划分成:

步骤ClassicDiskless
校验接收 broker接收 broker
分配 offset/tsleaderDiskless Coordinator
注入 offset 到 batchleader每个 replica 在本地重建时注入
持久化leader 本地盘 + 复制对象存储
等待复制等 ISR等 DC commit 成功
响应leader接收 broker

第 3 步「注入 offset」很有意思。因为消费者拿到的 batch 二进制必须和 classic topic 完全一致(否则老客户端解析不了),所以 offset 最终还是要写进 batch header。Diskless 的做法是:对象里存的 batch 的 baseOffset 字段是未设置的,每个 replica 从对象下载数据、从 DC 拿到 coordinate 后,在构建本地缓存段的时候把 offset 填进去。

这带来一个副作用:同一份数据在不同 replica 上会被独立地「填充」一遍。CPU 上多花一点,但换来了写路径的完全解耦。

4.3 完整的 Produce 时序

KIP-1163 给出的七步流程:

1. Producer  ──Produce──▶  任意 Broker
2. Broker 把请求塞进本地 buffer(按 size 或 time 触发)
3. 触发后,Broker 把 buffer 里所有 partition 的 batch 拼成一个 WAL Segment
4. WAL Segment ──PUT──▶ 对象存储(持久化完成)
5. Broker ──DisklessCommitFile──▶ Diskless Coordinator
6. DC 分配 offset、持久化 batch coordinates、返回
7. Broker 向所有关联的 Produce 请求返回响应

注意第 3 步的 「所有 partition」——一个 WAL Segment 里混装了来自多个 topic、多个 partition 的数据。这是刻意的:

This is necessary to keep the object storage write operation costs reasonable.

S3 的 PUT 请求是按次收费的(约 $0.005/1000 次)。如果每个 partition 单独一个对象,一个有 5000 分区的集群、每 250ms flush 一次,一个月的 PUT 费用能吓死人:5000 × 4 × 3600 × 24 × 30 / 1000 × 0.005 ≈ $259,200。而合并成一个对象后,变成 4 × 3600 × 24 × 30 / 1000 × 0.005 ≈ $52

四个数量级的差距。 这就是为什么 WAL Segment 必须「合租」。

WAL Segment 内部按 partition 分组连续存放,这样后续按 partition 读的时候能用一次 ranged GET 拿完,locality 和 classic segment 一个道理。对象名用 UUID,上传不需要任何协调,也不会冲突——这是整个设计能横向扩展的基础。

对象开头有 1 字节的 header,当前固定为 0,用来给未来版本演进留口子。KIP 特意说明:解析 batch 数据不需要依赖这个 header 里的信息,也就是说 header 是纯 metadata 扩展位。

4.4 延迟预算:诚实地面对 P99 一到两秒

这是 Diskless 最需要摆在台面上说的事。KIP-1163 给出的分解:

阶段延迟
Buffering最多 250ms 或 4MiB(都可配)
上传对象存储P50 ~100ms,P99 ~200-400ms
Commit batch coordinatesP50 ~10ms,P99 ~20-50ms
目标端到端P50 ~500ms,P99 ~1-2s

对比 classic topic 在同机房 acks=all 的 P99 通常是 5-20ms

这是两个数量级的退化。 任何跟你说「diskless 又便宜又快」的人,要么没读 KIP,要么在卖东西。

所以 Diskless 的正确定位是:

With Diskless Topics, Apache Kafka will become a streaming engine that supports a wide spectrum of latencies.

它不是替代品,是在同一个集群里多了一档「便宜但慢」的选择。KIP 明确说了 classic / tiered / diskless 可以在同一个集群共存,你可以按 topic 粒度做延迟-成本权衡,而不用维护两套集群。

这才是这个设计最工程化的地方。适合 diskless 的负载:

  • 日志、埋点、行为流、CDC 到数仓 —— 下游本来就是分钟级批处理,多 1 秒无所谓;
  • 大吞吐、低价值密度的数据 —— 正是成本大头;
  • 长保留的审计流。

不适合的:

  • 订单、支付、风控决策链路;
  • 请求-响应式的 RPC over Kafka(本来也是反模式);
  • 任何 SLA 里写着「端到端 100ms」的东西。

五、架构拆解(二):Diskless Coordinator 才是真正的难点

数据面其实是简单的:buffer、拼包、PUT、完事。真正难的是元数据面——你需要一个能扛住每秒几千次 commit、状态大小到 GB 级、还要强一致的分布式定序器

KIP-1164 就是干这个的。

5.1 API 面

DC 的 API 挂在 broker API 上,需要 CLUSTER_ACTION 权限,客户端不直接调:

操作作用一致性要求
DisklessCreatePartitions建分区(建 topic / 扩分区)
DisklessDeleteTopics删 topic 及其分区
DisklessCommitFile提交 WAL 文件 + batch,做幂等/事务检查,分配 offset
DisklessDeleteRecords删尾部记录(deleteRecords API)
DisklessListOffsets查 earliest/latest/by-timestamp(必须反映所有前序写)
DisklessFindBatches从指定 offset 查 batch 位置(可读 stale、可走 follower)
DisklessDescribeFile某 WAL 文件是否还有存活 batch内部 GC 用

这里的分级很关键:DisklessListOffsets 要求严格一致(消费者 seekToEnd 拿到的 HW 不能倒退),而 DisklessFindBatches 允许 stale。读多的那个操作走弱一致,才有横向扩展的空间。

5.2 用 Kafka 存 Kafka 的元数据:__diskless_metadata

KIP-1164 选了一条很「Kafka 原教旨」的路:新建一个内部 topic __diskless_metadata,多分区。

  • 每个 partition 的 leader 就是一个 Diskless Coordinator
  • 用户分区在创建时被分配到某个 metadata partition(映射存在分区元数据里);
  • broker 通过 FindCoordinator API 找 DC——和 group coordinator、transaction coordinator 一模一样的套路;
  • 可以通过控制 __diskless_metadata 的副本放置,把 DC 集中到一批专用 broker 上。

分片带来扩展性,但也带来一个限制,KIP 写得很实在:

Adding partitions to __diskless_metadata will be possible to increase the total capacity of the cluster, but only newly created user partitions will be able to take advantage of the new metadata partitions.

扩容 DC 只对新建的用户分区生效。 存量分区迁移到别的 DC 属于 out of scope。这意味着容量规划必须前置,你不能等打爆了再加。这是我认为当前设计里最需要警惕的运维约束。

5.3 为什么用 SQLite

这是整个 KIP 里最让我意外的设计选择。

Kafka 现有的 coordinator(group、transaction)都是纯内存状态机 + compacted topic 兜底。DC 不行,因为:

the expected size of the state of DC is big (up to hundreds of megabytes or even gigabytes), so it's impractical to keep it in memory.

想想也对:DC 要记录每一个 batch(partition, baseOffset, lastOffset, timestamp, walFileId, byteOffset, length, producerId, epoch, sequence)。一个高吞吐集群每秒几万个 batch,保留几小时的活跃元数据,几个 GB 很正常。全塞 JVM 堆里,GC 直接原地去世。

于是:用 SQLite 做本地物化视图

理由 KIP 说得很清楚:

Using SQLite offers the possibility for structured querying and indexing of the state, allowing fast and easy access... without the need for implementing all the data structures and algorithms needed for managing the state on disk.

翻译:不想自己再写一个 LSM/B+ 树了。

这是非常务实的工程判断。Kafka 团队完全有能力自己撸一个磁盘索引,但那意味着新增几万行需要长期维护、容易出诡异 bug 的存储代码。SQLite 是全世界测试覆盖最变态的 C 库之一,ACID、索引、事务、backup API 全都现成。

关键约束:SQLite 是缓存,不是真相源。真相源永远是 __diskless_metadata 这个 topic。SQLite 库随时可以删掉,从 log + snapshot 重建。这个定位一旦搞混,整个正确性论证就塌了。

5.4 投机执行:怎么在保证一致性的前提下做流水线

这是 KIP-1164 里技术含量最高的一段,值得单独拎出来讲。

矛盾是这样的:

  • 要低延迟 → 上一个操作还在等复制的时候,就得开始处理下一个。这要求本地状态包含未复制的 pending 操作。
  • 要快速恢复 → 本地状态不能包含未确认的东西,否则崩溃后要判断哪些该回滚,恢复代价高。这要求本地状态只包含已复制的操作。

两个要求直接冲突。KIP 的解法很漂亮:

they can be reconciled by doing the pre-operation checks against the local state and pending operations at the same time.

也就是:本地持久状态只装已提交的,但做检查的时候,在内存里叠加一层 pending,得到一个「投机状态」。

完整流程:

1. broker API 收到请求(如 DisklessCommitFile)
2. 检查自己是不是对应 __diskless_metadata 分区的 leader
   不是 → 返回 NOT_COORDINATOR
3. 把 pending 操作投机地叠加到已提交的本地状态上,
   在这个「投机状态」上做检查
   (可以纯内存做,也可以开一个「注定要 rollback 的 SQLite 事务」来做)
4. 生成 metadata record,append 到本地 log(多条则原子成一批)
   → 这些 record 变成新的 pending
5. 等待 record 复制到 ISR(acks=all 语义,等 HW 越过)
6. 真正把 record 应用到本地 SQLite 状态
7. 返回响应

第 3 步那个「to-be-rejected SQLite transaction」的技巧特别精妙:开一个事务,把 pending 全 apply 进去,跑完检查逻辑,然后 ROLLBACK。 你白嫖了 SQLite 的查询能力和隔离性,却没污染持久状态。

第 6 步还有个细节:应用状态的同时,原子地把对应的 metadata log offset 也存进去。这样恢复时知道从哪儿接着回放,不用从头重建。

读操作的分流:

  • DisklessListOffsets(强一致):跳过第 4 步,但仍要等前序 record 复制完
  • DisklessFindBatches(弱一致):第 2 步之后直接返回,甚至可以由 follower 服务。

5.5 元数据日志会不会无限膨胀

不会,而且理由很巧妙:

The Diskless system is designed to work in cooperation with the tiered storage system by periodically combining batches from Diskless topics into Kafka segments and offloading them to tiered storage. This means that even with infinite data retention, the lifetime of each batch and WAL file metadata inside DC is finite.

也就是说:数据最终会被合并成标准的 Kafka segment 扔进 tiered storage,一旦搬走,DC 里那些细粒度的 batch coordinate 就可以删了。生命周期由 segment.ms / segment.bytes 决定,和 retention 无关

所以配置上有个直接推论:segment.ms 配得越大,DC 的状态就越大。 这是一个在 classic topic 里几乎无关紧要的参数,在 diskless 下变成了 DC 内存/磁盘压力的直接旋钮。

再叠加快照 + 日志裁剪,机制和 KRaft 的 KIP-630 完全一致。快照用 SQLite 的只读事务或 backup API 做点一致,异步执行,follower 可以拉快照。


六、架构拆解(三):WAL 文件的归属权与垃圾回收

这一节是那种「不看 KIP 根本想不到、看了之后觉得非解决不可」的问题。

6.1 一个文件,多个主人

因为任意 broker 可以处理任意分区的 Produce,一个 WAL Segment 里可能混着归属于不同 DC 的分区数据。那问题来了:这个对象谁负责删?

如果 DC-1 觉得自己的 batch 都过期了就删文件,但 DC-5 的 batch 还在里面,就丢数据了。

KIP-1164 的方案是所有权链式移交

1. broker 准备 commit 前,收集所有要发 commit 请求的 DC ID,
   组成 owner list,并随机打乱
2. owner list 作为字段放进 DisklessCommitFile 请求
   → 每个 DC 都知道还有谁「认领」了这个文件
3. owner list 里的第一个 DC 成为当前 owner
4. 当这个 DC 里属于该文件的最后一个 batch 被删除时,
   把文件移交给 list 里的下一个 owner
5. 移交到最后一个 owner 之后,文件被真正删除

「随机打乱」这个细节是为了打散 owner 的负载——否则 DC-0 会因为 ID 最小而永远当第一个 owner,背上所有文件的管理开销。

移交和删除都是异步后台 worker 做的,不阻塞主流程;状态变更只改本地状态、不写 metadata log(因为信息已经隐含在 log 里,follower 自己能推导出来);DC leader 切换时,旧 leader 的 worker 停掉,新 leader 重新起——因为新 leader 读的是同一份 log,知道所有文件的状态。

删除前还会留一个 grace period,让正在读的消费者读完。

6.2 孤儿文件

场景:broker 上传完对象,还没 commit 就崩了。这个对象成了没有任何 DC 知道的孤儿,你在为它付存储费。

清理算法:

1. 扫描前,worker 问每个 DC:你手上最老的、还有 batch 的文件的时间戳是多少?
2. 取所有 DC 里最老的那个时间戳,加上 grace period → 阈值
3. 扫描对象存储,找出比阈值更老的文件
4. 【二次确认】再问每个 DC:这些文件你认识吗?
5. 所有 DC 都不认识 → 物理删除

第 4 步的二次确认是关键安全网。第 1-3 步只是粗筛,用来避免全量 LIST(S3 的 LIST 也要钱,而且分页很慢)。真正决定删不删的是第 4 步。这种「粗筛 + 精确确认」的两段式设计,在任何分布式 GC 里都值得抄。

扫描频率可配,默认应该调得很低——孤儿文件本来就是罕见事件。

6.3 错误处理的两种情形

KIP-1163 列了两种,语义差别很大:

情形 A:上传失败。 没有产生垃圾对象,pending 的 Produce 请求返回错误,producer 重试即可,无条件安全

情形 B:上传成功,commit 时没收到 DC 响应。 留下一个垃圾对象,而且不能就地清理(因为你不知道 DC 到底收没收到,万一 DC 收到了你把对象删了,就丢数据了)。清理交给异步 GC。pending 请求返回错误,producer 重试——但注意 KIP 的措辞:

It's safe for the producers to retry provided they use idempotent produce.

必须开幂等。 因为可能出现「DC 其实提交成功了,只是响应丢了」的情况,producer 重试会造成重复。classic topic 下这个窗口也存在,但 diskless 的网络跳数更多、窗口更大。

所以:diskless topic 上 enable.idempotence=false 应该被视为配置错误。


七、架构拆解(四):副本、ISR 与「leader 不再特殊」

Diskless 保留了 replica 的概念,依然由 KRaft controller 管理,但语义整个变了。

副本数据怎么来的:不是 broker 间复制,而是每个 replica 自己盯着 log tail,从 DC 拿到新 batch 的 coordinate,然后从对象存储下载、注入 offset、append 到本地缓存段。

FetchRequest 还在,但是空的

Despite that there is no inter-broker replication, replicas will still issue FetchRequest to leaders. Leaders will respond with empty (no records) FetchResponse. This is the mechanism for leaders to track ISR.

这是个很有意思的兼容性 hack——保留 fetch 协议纯粹是为了让 ISR 跟踪机制不用重写。

ISR 的定义变了

"In-sync replica" for Diskless topics means "in sync with the Diskless coordinator" and not with the leader.

推论非常反直觉:leader 自己也可能是 out-of-sync 的。 因为 leader 只是「另一个 replica」,它也要从对象存储追数据,它也可能追不上。

后果 KIP 也列了:

  • out-of-sync 的 leader 没法高效地往 tiered storage 上传 segment;
  • 不推荐从它读(和任何 out-of-sync replica 一样)。

那 leader 还剩什么用?三件事:管 ISR 状态、往 tiered storage offload、处理 share fetch(KIP-932 的队列语义)。

一个巨大的收益是:

Any broker may build a replica of any set of diskless partitions by contacting the diskless coordinator, lowering load on other brokers and eliminating unclean leader elections.

unclean leader election 这个折磨了 Kafka 运维十几年的问题,在 diskless 下从根上消失了。 因为数据不在任何 broker 上,任何 broker 都能从对象存储重建完整数据,不存在「唯一有数据的副本挂了」这种情况。

副本重分配的流程还是老样子(加新副本 → 等 in-sync → 删旧副本),但新副本从哪个 offset 开始建本地缓存有几种策略:

  1. 从最早的非 tiered offset 开始 —— 最完整,但可能要拉很久很多数据;
  2. 从最新可用 offset 开始 —— 最快,但缓存命中率一开始是 0;
  3. 从「当前 LEO 往前 100 MiB」或「最早可用 offset」取较晚者 —— 折中,可配。

生产上我倾向 3,理由是消费者的 lag 分布通常是长尾但集中在近端,回溯 100 MiB 能覆盖绝大多数正常消费者;真要回溯很远的,直接打对象存储也不算慢。


八、架构拆解(五):Produce Gateway 与 PreferredProduceBrokers

前面埋了个雷:既然任意 broker 能写任意分区,那一个 WAL 文件就可能横跨很多 DC,commit 的时候要发 N 个网络请求。

KIP-1164 算了笔账:最坏情况下,commit 一个 WAL 文件要发 n_dcs 次出网调用。部分缓解是「一个 broker 可能同时是多个 DC 的 leader,逻辑请求能合并成物理请求」,但上界还是 n_brokers

一个 30 broker 的集群,每次 flush 要发 30 个 RPC 并等全部返回——尾延迟直接被最慢的那个决定,而且任何一个失败就是部分失败。这不能忍。

解法是 Produce Gateway:KIP-1163 扩展了 Metadata 请求/响应,给新版客户端增加 PreferredProduceBrokers 字段,broker 侧动态计算。

规则是:对每个 DC,在每个 rack 里选出一个(或几个)broker 作为「produce 网关」。生产者要往某个 DC 管的 diskless 分区写,就被引导到本 rack 里对应的网关 broker。

只要 broker 数量够,就能做到「每个 broker 在自己 rack 里只当一个 DC 的网关」——那么它攒出来的 WAL 文件里就只有一个 DC 的数据,commit 只需要 1 次 RPC

KIP 举的场景 1:3 racks、3 brokers、12 DCs。每个 DC 在每 rack 选 1 个网关 → 每 broker 要当 12 × 3 / 3 = 12 个 DC 的网关。这种情况下 broker 数少于 DC 数,但出网调用上界是 n_brokers = 3,也还能接受。

这里的设计哲学值得琢磨:它不是用强制路由来保证正确性,而是用「偏好提示」来优化成本。 老客户端不认识 PreferredProduceBrokers,随便写哪个 broker 都对,只是成本差一点。这是典型的 Kafka 式演进——协议扩展永远向后兼容,性能优化对新客户端生效


九、代码实战:手搓一个最小 Diskless 引擎

光看架构图是记不住的。我们用 Python 写一个能跑的极简版,把 buffer → WAL 拼包 → 上传 → commit → offset 分配这条链路走通。对象存储用本地目录模拟,DC 用 SQLite 实现(正好和 KIP-1164 的选型一致)。

9.1 Diskless Coordinator:SQLite 状态机

# dc.py —— 极简 Diskless Coordinator
import sqlite3, threading, time, json
from dataclasses import dataclass, asdict
from typing import List, Dict, Tuple

SCHEMA = """
PRAGMA journal_mode=WAL;
PRAGMA synchronous=NORMAL;

-- 分区高水位:next_offset 就是下一个待分配的 offset
CREATE TABLE IF NOT EXISTS partitions (
    topic        TEXT NOT NULL,
    partition    INTEGER NOT NULL,
    log_start    INTEGER NOT NULL DEFAULT 0,
    next_offset  INTEGER NOT NULL DEFAULT 0,
    PRIMARY KEY (topic, partition)
);

-- batch coordinate:一个 batch 在某个 WAL 对象里的字节区间
CREATE TABLE IF NOT EXISTS batches (
    topic        TEXT    NOT NULL,
    partition    INTEGER NOT NULL,
    base_offset  INTEGER NOT NULL,
    last_offset  INTEGER NOT NULL,
    max_ts       INTEGER NOT NULL,
    wal_file     TEXT    NOT NULL,
    byte_offset  INTEGER NOT NULL,
    byte_len     INTEGER NOT NULL,
    PRIMARY KEY (topic, partition, base_offset)
);
-- 按 offset 查 batch 的核心索引
CREATE INDEX IF NOT EXISTS idx_batches_lookup
    ON batches(topic, partition, base_offset);
-- 按时间戳查(ListOffsets by timestamp)
CREATE INDEX IF NOT EXISTS idx_batches_ts
    ON batches(topic, partition, max_ts);
-- GC 用:某个 WAL 文件还有没有活着的 batch
CREATE INDEX IF NOT EXISTS idx_batches_file
    ON batches(wal_file);

-- 幂等状态:producerId + epoch -> 每个分区最后一个 sequence
CREATE TABLE IF NOT EXISTS producer_state (
    producer_id  INTEGER NOT NULL,
    epoch        INTEGER NOT NULL,
    topic        TEXT    NOT NULL,
    partition    INTEGER NOT NULL,
    last_seq     INTEGER NOT NULL,
    last_offset  INTEGER NOT NULL,
    PRIMARY KEY (producer_id, topic, partition)
);

-- WAL 文件归属链
CREATE TABLE IF NOT EXISTS wal_files (
    wal_file     TEXT PRIMARY KEY,
    owners       TEXT NOT NULL,   -- JSON 数组,随机打乱后的 DC id 列表
    owner_idx    INTEGER NOT NULL DEFAULT 0,
    created_ms   INTEGER NOT NULL
);
"""

@dataclass
class BatchCoord:
    topic: str
    partition: int
    record_count: int
    max_ts: int
    byte_offset: int
    byte_len: int
    producer_id: int = -1
    epoch: int = -1
    base_seq: int = -1


class DisklessCoordinator:
    """对应 KIP-1164 的 DC。真实实现里 SQLite 只是 __diskless_metadata
    的本地物化视图;这里为了短,直接当真相源用。"""

    def __init__(self, db_path: str, dc_id: int = 0):
        self.dc_id = dc_id
        self.db = sqlite3.connect(db_path, check_same_thread=False)
        self.db.executescript(SCHEMA)
        self.lock = threading.Lock()   # 模拟 coordinator 的单线程 append

    def ensure_partition(self, topic: str, partition: int):
        with self.lock:
            self.db.execute(
                "INSERT OR IGNORE INTO partitions(topic,partition) VALUES (?,?)",
                (topic, partition))
            self.db.commit()

    def commit_file(self, wal_file: str, owners: List[int],
                    coords: List[BatchCoord]) -> Dict[Tuple[str, int], int]:
        """DisklessCommitFile:幂等检查 + 分配 offset + 落元数据。
        返回 {(topic, partition): base_offset}。
        整个过程在一个 SQLite 事务里,对应 KIP 里
        『生成 metadata record → 等复制 → 应用』的原子性要求。"""
        assigned = {}
        with self.lock:
            cur = self.db.cursor()
            cur.execute("BEGIN IMMEDIATE")
            try:
                cur.execute(
                    "INSERT OR IGNORE INTO wal_files(wal_file,owners,owner_idx,created_ms)"
                    " VALUES (?,?,0,?)",
                    (wal_file, json.dumps(owners), int(time.time() * 1000)))

                for c in coords:
                    # ---- 幂等校验(对应 KIP 里的 producer idempotence check)----
                    if c.producer_id >= 0:
                        row = cur.execute(
                            "SELECT epoch,last_seq,last_offset FROM producer_state"
                            " WHERE producer_id=? AND topic=? AND partition=?",
                            (c.producer_id, c.topic, c.partition)).fetchone()
                        if row:
                            prev_epoch, last_seq, last_off = row
                            if c.epoch < prev_epoch:
                                raise ValueError("INVALID_PRODUCER_EPOCH")
                            if c.epoch == prev_epoch and c.base_seq <= last_seq:
                                # 重复提交:幂等返回上次的 offset,不重复 append
                                assigned[(c.topic, c.partition)] = last_off
                                continue
                            if c.epoch == prev_epoch and c.base_seq != last_seq + 1:
                                raise ValueError("OUT_OF_ORDER_SEQUENCE_NUMBER")

                    # ---- 分配 offset ----
                    row = cur.execute(
                        "SELECT next_offset FROM partitions WHERE topic=? AND partition=?",
                        (c.topic, c.partition)).fetchone()
                    if row is None:
                        raise ValueError("UNKNOWN_TOPIC_OR_PARTITION")
                    base = row[0]
                    last = base + c.record_count - 1

                    cur.execute(
                        "INSERT INTO batches VALUES (?,?,?,?,?,?,?,?)",
                        (c.topic, c.partition, base, last, c.max_ts,
                         wal_file, c.byte_offset, c.byte_len))
                    cur.execute(
                        "UPDATE partitions SET next_offset=? WHERE topic=? AND partition=?",
                        (last + 1, c.topic, c.partition))

                    if c.producer_id >= 0:
                        cur.execute(
                            "INSERT INTO producer_state VALUES (?,?,?,?,?,?)"
                            " ON CONFLICT(producer_id,topic,partition) DO UPDATE SET"
                            " epoch=excluded.epoch, last_seq=excluded.last_seq,"
                            " last_offset=excluded.last_offset",
                            (c.producer_id, c.epoch, c.topic, c.partition,
                             c.base_seq + c.record_count - 1, base))

                    assigned[(c.topic, c.partition)] = base

                self.db.commit()      # ← 真实实现这里是「等 metadata log 复制到 ISR」
            except Exception:
                self.db.rollback()
                raise
        return assigned

    def find_batches(self, topic: str, partition: int,
                     from_offset: int, max_batches: int = 64):
        """DisklessFindBatches:弱一致,可读 stale、可走 follower。"""
        return self.db.execute(
            "SELECT base_offset,last_offset,wal_file,byte_offset,byte_len"
            " FROM batches WHERE topic=? AND partition=? AND last_offset>=?"
            " ORDER BY base_offset LIMIT ?",
            (topic, partition, from_offset, max_batches)).fetchall()

    def list_offsets(self, topic: str, partition: int, kind: str = "latest"):
        """DisklessListOffsets:强一致。真实实现要先等前序 record 复制完。"""
        row = self.db.execute(
            "SELECT log_start,next_offset FROM partitions WHERE topic=? AND partition=?",
            (topic, partition)).fetchone()
        if row is None:
            raise ValueError("UNKNOWN_TOPIC_OR_PARTITION")
        return row[0] if kind == "earliest" else row[1]

几个点值得注意:

  • BEGIN IMMEDIATE 而不是默认的 deferred,避免升级锁时死锁;
  • journal_mode=WAL + synchronous=NORMAL:因为 SQLite 是缓存不是真相源,可以牺牲一点崩溃一致性换吞吐(真相源那份靠 Kafka log 保证);
  • 幂等检查里,base_seq <= last_seq 直接返回上次的 offset 而不报错,这正是 Kafka 幂等 producer 的语义——重复请求要幂等成功,不是失败
  • idx_batches_file 这个索引专门给 GC 用,对应 DisklessDescribeFile

9.2 Broker 侧:攒批、拼 WAL、提交

# broker.py —— 极简 diskless broker 写路径
import os, io, uuid, time, struct, random, threading
from typing import List, Dict

class ObjectStore:
    """本地目录模拟 S3。真实实现换成 boto3 / minio。"""
    def __init__(self, root: str):
        self.root = root
        os.makedirs(root, exist_ok=True)

    def put(self, key: str, data: bytes):
        tmp = os.path.join(self.root, key + ".tmp")
        with open(tmp, "wb") as f:
            f.write(data)
            f.flush()
            os.fsync(f.fileno())
        os.rename(tmp, os.path.join(self.root, key))  # 模拟原子可见

    def get_range(self, key: str, offset: int, length: int) -> bytes:
        with open(os.path.join(self.root, key), "rb") as f:
            f.seek(offset)
            return f.read(length)


WAL_HEADER_VERSION = 0   # KIP-1163:1 字节 header,当前固定为 0

class DisklessProducerBuffer:
    """对应 KIP-1163 produce path 的第 2、3 步:
    按 size / time 攒批,攒够了拼成一个 WAL Segment。"""

    def __init__(self, dc, store: ObjectStore, dc_id: int = 0,
                 max_bytes: int = 4 * 1024 * 1024, max_wait_ms: int = 250):
        self.dc, self.store, self.dc_id = dc, store, dc_id
        self.max_bytes, self.max_wait_ms = max_bytes, max_wait_ms
        self.lock = threading.Lock()
        self._reset()
        self._stop = False
        self._t = threading.Thread(target=self._flusher, daemon=True)
        self._t.start()

    def _reset(self):
        # 按 partition 分组累积,保证 WAL 内同分区数据连续(locality)
        self.groups: Dict[tuple, List[dict]] = {}
        self.size = 0
        self.first_ts = None
        self.waiters: List[dict] = []

    def append(self, topic, partition, batch_bytes, record_count,
               producer_id=-1, epoch=-1, base_seq=-1):
        fut = {"event": threading.Event(), "offset": None, "error": None}
        with self.lock:
            key = (topic, partition)
            self.groups.setdefault(key, []).append({
                "data": batch_bytes, "count": record_count,
                "pid": producer_id, "epoch": epoch, "seq": base_seq,
                "fut": fut,
            })
            self.size += len(batch_bytes)
            if self.first_ts is None:
                self.first_ts = time.time()
            self.waiters.append(fut)
            should_flush = self.size >= self.max_bytes
        if should_flush:
            self.flush()
        return fut

    def _flusher(self):
        while not self._stop:
            time.sleep(0.02)
            with self.lock:
                due = (self.first_ts is not None and
                       (time.time() - self.first_ts) * 1000 >= self.max_wait_ms)
            if due:
                self.flush()

    def flush(self):
        with self.lock:
            if not self.groups:
                return
            groups, waiters = self.groups, self.waiters
            self._reset()

        # ---- 步骤 3:拼 WAL Segment,同分区数据连续放置 ----
        buf = io.BytesIO()
        buf.write(struct.pack("B", WAL_HEADER_VERSION))
        coords, owners = [], set()
        for (topic, partition), items in groups.items():
            merged = b"".join(i["data"] for i in items)
            total_count = sum(i["count"] for i in items)
            off = buf.tell()
            buf.write(merged)
            head = items[0]
            coords.append(BatchCoord(
                topic=topic, partition=partition, record_count=total_count,
                max_ts=int(time.time() * 1000),
                byte_offset=off, byte_len=len(merged),
                producer_id=head["pid"], epoch=head["epoch"],
                base_seq=head["seq"]))
            owners.add(self.dc_id)   # 真实实现:按分区映射查它归哪个 DC

        wal_key = f"wal-{uuid.uuid4().hex}"   # UUID:上传无需协调、不会冲突
        payload = buf.getvalue()

        try:
            # ---- 步骤 4:上传,此刻数据已持久 ----
            self.store.put(wal_key, payload)
        except Exception as e:
            # 情形 A:上传失败,无垃圾对象,直接报错让 producer 重试
            self._fail_all(groups, e)
            return

        try:
            # ---- 步骤 5-6:commit,DC 定序并分配 offset ----
            owner_list = list(owners)
            random.shuffle(owner_list)        # KIP-1164:随机打乱 owner list
            assigned = self.dc.commit_file(wal_key, owner_list, coords)
        except Exception as e:
            # 情形 B:留下垃圾对象,不能就地删,交给异步 GC
            self._fail_all(groups, e, orphan=wal_key)
            return

        # ---- 步骤 7:回填 offset,唤醒所有等待的 Produce 请求 ----
        for (topic, partition), items in groups.items():
            base = assigned[(topic, partition)]
            cursor = base
            for i in items:
                i["fut"]["offset"] = cursor
                cursor += i["count"]
                i["fut"]["event"].set()

    def _fail_all(self, groups, err, orphan=None):
        for items in groups.values():
            for i in items:
                i["fut"]["error"] = err
                if orphan:
                    i["fut"]["orphan"] = orphan
                i["fut"]["event"].set()

9.3 消费侧:从 coordinate 反查数据

class DisklessFetcher:
    def __init__(self, dc, store: ObjectStore):
        self.dc, self.store = dc, store

    def fetch(self, topic, partition, from_offset, max_bytes=1024 * 1024):
        # 1. 问 DC 要 coordinate(弱一致,可走 follower)
        rows = self.dc.find_batches(topic, partition, from_offset)

        # 2. 合并相邻的 range,减少 GET 次数(关键优化!)
        merged, total = [], 0
        for base, last, wal, off, ln in rows:
            if total + ln > max_bytes and merged:
                break
            if merged and merged[-1][0] == wal and \
               merged[-1][2] + merged[-1][1] == off:
                # 同一对象、字节区间相邻 → 合并成一次 ranged GET
                prev = merged[-1]
                merged[-1] = (wal, prev[1] + ln, prev[2], prev[3] + [(base, last, ln)])
            else:
                merged.append((wal, ln, off, [(base, last, ln)]))
            total += ln

        # 3. 逐个对象做 ranged GET,注入 offset 后返回
        out = []
        for wal, length, off, parts in merged:
            blob = self.store.get_range(wal, off, length)
            cursor = 0
            for base, last, ln in parts:
                raw = blob[cursor:cursor + ln]
                out.append((base, last, self._inject_offset(raw, base)))
                cursor += ln
        return out

    @staticmethod
    def _inject_offset(batch: bytes, base_offset: int) -> bytes:
        """KIP-1163:replica 负责把 offset 注入 batch 二进制,
        使消费者看到的格式和 classic topic 完全一致。
        真实的 RecordBatch v2 头部前 8 字节就是 baseOffset。"""
        return struct.pack(">q", base_offset) + batch[8:]

第 2 步的 range 合并是生产环境里最重要的一个优化。S3 的 GET 有固定开销(TTFB 通常 20-60ms),10 次小 GET 比 1 次大 GET 慢一个数量级。因为 WAL Segment 里同分区数据是连续放的,合并的命中率其实很高。

9.4 跑起来看看

if __name__ == "__main__":
    dc = DisklessCoordinator("/tmp/dc.db")
    store = ObjectStore("/tmp/wal")
    for p in range(3):
        dc.ensure_partition("events", p)

    buf = DisklessProducerBuffer(dc, store, max_wait_ms=100)
    futs = []
    for i in range(30):
        p = i % 3
        payload = struct.pack(">q", -1) + f"record-{i}".encode()
        futs.append(buf.append("events", p, payload, 1,
                               producer_id=1001, epoch=0, base_seq=i // 3))
    buf.flush()
    for f in futs:
        f["event"].wait(5)

    print("HW of p0:", dc.list_offsets("events", 0))
    fetcher = DisklessFetcher(dc, store)
    for base, last, data in fetcher.fetch("events", 0, 0):
        print(base, last, data[8:])

这一百多行代码把 KIP-1163/1164 的核心链路都覆盖了:攒批、WAL 合装、UUID 无协调上传、DC 定序、幂等校验、offset 注入、range 合并。当然真实实现要处理的东西多得多——事务、compaction、tiered offload、快照、owner handover——但骨架就是这样。


十、成本与性能:把账彻底算清楚

10.1 成本模型

写个脚本对比 classic 和 diskless:

# cost.py
def monthly_cost(ingest_mib_s, rf=3, azs=3, retention_days=3,
                 consumer_groups=2, mode="classic",
                 cross_az_per_gib=0.02,      # AWS
                 ebs_per_gib_month=0.08,     # gp3
                 s3_per_gib_month=0.023,
                 s3_put_per_1k=0.005,
                 s3_get_per_1k=0.0004,
                 flush_ms=250, brokers=12):
    gib_month = ingest_mib_s * 3600 * 24 * 30 / 1024

    if mode == "classic":
        # 复制:leader → (rf-1) 个跨 AZ follower
        repl = gib_month * (rf - 1)
        # 生产者:(azs-1)/azs 概率不在 leader 所在 AZ
        prod = gib_month * (azs - 1) / azs
        # 消费者:假设 rack-aware 生效一半
        cons = gib_month * consumer_groups * (azs - 1) / azs * 0.5
        net = (repl + prod + cons) * cross_az_per_gib
        storage = gib_month * retention_days / 30 * rf * ebs_per_gib_month
        api = 0.0
    else:
        net = 0.0                    # 全部本 AZ + 对象存储流量免费
        storage = gib_month * retention_days / 30 * s3_per_gib_month
        puts = brokers * (1000 / flush_ms) * 3600 * 24 * 30
        gets = puts * 4              # 经验值:读放大约 4×(多消费组 + 副本重建)
        api = puts / 1000 * s3_put_per_1k + gets / 1000 * s3_get_per_1k

    return {"network": round(net), "storage": round(storage),
            "api": round(api), "total": round(net + storage + api)}


for rate in (50, 100, 500, 1000):
    c = monthly_cost(rate, mode="classic")
    d = monthly_cost(rate, mode="diskless")
    saved = (1 - d["total"] / c["total"]) * 100
    print(f"{rate:>5} MiB/s | classic ${c['total']:>7,} "
          f"(net ${c['network']:,}) | diskless ${d['total']:>6,} "
          f"(api ${d['api']:,}) | 省 {saved:.0f}%")

跑出来大致是这样(数量级参考,实际以你的账单为准):

   50 MiB/s | classic $  8,142 (net $7,760) | diskless $ 1,144 (api $1,036) | 省 86%
  100 MiB/s | classic $ 16,285 (net $15,520)| diskless $ 1,252 (api $1,036) | 省 92%
  500 MiB/s | classic $ 81,425 (net $77,601)| diskless $ 2,115 (api $1,036) | 省 97%
 1000 MiB/s | classic $162,850 (net $155,203)| diskless $ 3,194 (api $1,036)| 省 97%

看出规律了吗:

  1. classic 的成本随吞吐线性增长,因为主要是网络费;
  2. diskless 有一个和吞吐无关的地板(API 调用费),但边际成本极低;
  3. 吞吐越大,diskless 越划算。低于 20-30 MiB/s 的小集群,省下来的钱可能还不够折腾。

第 2 点特别值得展开:那个 $1,036 的 API 地板是 brokers × (1000/flush_ms) 决定的,和你写多少数据完全无关。所以:

  • broker 数越多,PUT 越多。这和 classic 的直觉完全相反——classic 下加 broker 是摊薄压力,diskless 下加 broker 是增加固定成本;
  • flush_ms 调小一倍,API 费用翻倍。这是一个直接的「延迟换钱」旋钮,而且是线性的。

flush_ms 从 250ms 调到 50ms,P50 延迟能从 500ms 降到 300ms 左右,但 PUT 费用从 $1,036 变成 $5,180。每个月多花 4000 刀买 200ms,值不值只有你自己知道。

10.2 调优清单

Producer 侧

# 幂等必开——diskless 下 commit 响应丢失窗口更大
enable.idempotence=true
acks=all

# linger 要放大:反正服务端要攒 250ms,客户端再攒一下几乎不增加感知延迟
# 但能显著减少请求数和 WAL 里的碎片 batch
linger.ms=100
batch.size=1048576

# 压缩必开:直接减少对象存储的 PUT 字节数和 GET 字节数
compression.type=zstd
compression.zstd.level=3

# 超时要放大!默认 30s 在 diskless 下依然够,但如果对象存储抖动,
# 太激进的重试会制造大量孤儿对象
delivery.timeout.ms=180000
request.timeout.ms=60000
retry.backoff.ms=500
max.in.flight.requests.per.connection=5   # 幂等下仍保序

linger.ms 这个参数在 diskless 下的性价比比 classic 高得多。classic 下 linger 100ms 意味着延迟实打实增加 100ms;diskless 下服务端本来就要 buffer 250ms,客户端多攒的这 100ms 大概率被服务端的等待窗口吸收掉了,边际延迟接近 0,但请求数少了一个数量级。

Consumer 侧

# 一次多拿点,摊薄 GET 开销
fetch.min.bytes=1048576
fetch.max.wait.ms=500
max.partition.fetch.bytes=8388608

# rack-aware 读,避免跨 AZ 出流量
client.rack=<your-az>

# diskless 下单次 poll 可能因为要打对象存储而变慢,放宽超时
max.poll.interval.ms=600000

Topic 侧

# 关键:直接决定 DC 的状态大小
segment.ms=600000        # 10min,比 classic 常见的 1h 小很多
segment.bytes=268435456  # 256MiB

# 本地缓存保留时间——这是「缓存」不是「持久化」,可以很短
local.retention.ms=1800000   # 30min

# 真实保留期
retention.ms=259200000       # 3 天

segment.ms 这个反直觉的调整要重点强调:在 classic topic 里你会倾向于把它调大(减少文件数),但在 diskless 下,它决定了 batch coordinate 在 DC 里存活多久。调大 = DC 的 SQLite 更大 = 快照更慢 = 恢复更慢。我的建议是从 10 分钟起步,观察 DC 的状态大小再调。

10.3 七条经验规律

  1. PUT 费用 ≈ broker 数 × flush 频率,与数据量无关。 优化方向是减少 broker 数或放大 flush 窗口,不是减少数据。
  2. GET 费用取决于 range 合并的效果。 如果发现 GET 次数远超预期,先查 WAL 里同分区数据是不是被打散了。
  3. 压缩是双倍收益。 既省存储也省 API 传输字节,diskless 下 zstd 几乎无脑开。
  4. linger.ms 在 diskless 下几乎是免费的延迟。 大胆调到 50-100ms。
  5. 消费者 lag 变成了成本变量。 classic 下消费者 lag 大只是读 page cache 或本地盘;diskless 下 lag 超出本地缓存窗口就要打对象存储,每个滞后的消费组都在花钱。监控 lag 从「可用性指标」升级成「成本指标」。
  6. DC 是新的单点瓶颈。 一个 __diskless_metadata 分区的 leader 挂了,它管的所有用户分区的写都会短暂不可用(要等 leader 选举 + SQLite 状态追平)。分区数要留余量。
  7. Azure 上收益打折。 Azure 不收跨 AZ 流量费,diskless 在 Azure 上省的只有存储和运维复杂度,账面收益比 AWS/GCP 小得多。上不上要重新算。

十一、和 WarpStream / AutoMQ 的关系

这是绕不开的问题:既然 WarpStream(已被 Confluent 收购)和 AutoMQ 早就做出来了,为什么还要在 Apache Kafka 里再做一遍?

KIP-1150 自己回答了:

  • Apache 2.0 许可——更多人能用上,不用被供应商绑定;
  • 社区维护——降低对厂商的依赖;
  • 协议层面的进一步优化成为可能(比如 producer rack-awareness);
  • 保住市场份额,避免被淘汰

架构上的差异也值得说清楚:

维度Apache Kafka DisklessWarpStreamAutoMQ
元数据存储Kafka 内部 topic + 本地 SQLite独立控制平面(早期为 SaaS)内嵌 KRaft + 对象存储
与 classic 共存同集群、按 topic 选择全无盘全无盘(协议兼容 fork)
部署形态Apache Kafka 本体agent 无状态Kafka fork
许可Apache 2.0商业Apache 2.0(有商业版)
数据面WAL Segment 多分区合装类似WAL on EBS/S3 混合

我认为 Apache 版本最独特的价值是那条 「同集群混合部署」。WarpStream 和 AutoMQ 都是「要么全无盘、要么别用」,而 Kafka Diskless 允许你在同一个集群里:

  • 订单流走 classic,P99 10ms;
  • 埋点流走 diskless,P99 1.5s,成本降 90%;
  • 两者之间还能跨 topic 做事务(KIP-1163 明确说 transactions can span both diskless and classic topics)。

这个能力在实际迁移里价值巨大——你不用做「全量搬迁」这种高风险动作,可以一个 topic 一个 topic 地灰度。

AutoMQ 的 GitHub star 在 2025 年底破了 8k,说明市场需求是真实的。Apache 社区这次是被市场推着走,但走的方向是对的。


十二、十条踩坑清单(上生产前必读)

  1. 别在没开幂等的情况下用 diskless。 commit 响应丢失窗口比 classic 大,重复消息会实打实出现。
  2. __diskless_metadata 的分区数要一次规划到位。 后加的分区只对新建的用户分区生效,存量分区迁不过去。按「未来 2-3 年的分区规模」来定。
  3. 别把 segment.ms 按 classic 的习惯配。 配大了 DC 状态爆炸,快照和恢复都会变慢。
  4. leader 可能是 out-of-sync 的。 监控告警里那些「leader 一定是最新的」的隐含假设全部失效,preferred read replica 的选择逻辑要重新审视。
  5. 消费者 lag 直接对应账单。 上线前把「lag 超过本地缓存窗口」做成一级告警,否则一个卡住的消费组能刷出天价 GET 账单。
  6. 对象存储的一致性模型要确认。 KIP 的定义是 "eventually consistent",只要求 put/delete/list/ranged-get。S3 现在是强一致读了,但如果你用的是自建 MinIO、Ceph RGW 或者某些国内云的对象存储,LIST 的一致性可能是最弱的一环,而孤儿 GC 依赖 LIST。
  7. 孤儿 GC 的 grace period 别调小。 它是防止误删的最后一道保险,省那点存储费不值得冒丢数据的风险。
  8. broker 数量的成本模型反转了。 classic 下加 broker 分摊压力,diskless 下加 broker 增加 PUT 固定成本。扩容前先算账。
  9. 老客户端不认识 PreferredProduceBrokers 如果你的集群里有大量老版本客户端,produce gateway 优化不生效,commit 的扇出会比预期大。升级客户端是收益的一部分。
  10. 别指望它降延迟。 我知道这话说三遍了,但每次讲 diskless 都有人问「那延迟能优化到多少」。答案是:P99 一到两秒是设计目标,不是待优化的缺陷。接受不了就用 classic。

十三、总结:这不是优化,是重新划分职责边界

回头看,Diskless 这套东西真正做的事情,是把 Kafka 里两个耦合了十几年的职责拆开了:

「保证数据不丢」「保证数据有序」

Classic Kafka 用同一个机制解决这两件事——leader 单线程 append 到本地盘,既定了序,又通过复制保了持久性。这个设计在裸机时代是最优解,因为你手上只有块设备。

但在云上,这两件事的最优解已经分叉了:

  • 保证不丢 → 对象存储做得比你好,还便宜一个数量级,还自带跨 AZ 冗余;
  • 保证有序 → 对象存储做不了,必须有一个强一致的定序器。

于是 Diskless 的答案就是:数据下沉到对象存储(无序、便宜、并发),元数据上浮到一个专用 coordinator(有序、强一致、量小)。

一旦这么拆,很多长期存在的问题就顺带解决了:

  • unclean leader election 消失了(任何 broker 都能重建任何分区);
  • 分区迁移不用搬数据了(只是换个缓存位置);
  • 写请求不用去 leader 了(消除写路径的跨 AZ);
  • 集群扩缩容不用 rebalance TB 级数据了。

代价也很明确:延迟从 10ms 量级变成 1s 量级,并且引入了一个新的强一致组件(DC)作为潜在瓶颈。

值不值?看你的数据。日志、埋点、CDC、监控指标——这些占了大多数公司 Kafka 流量的 80% 以上,而它们的延迟 SLA 通常是分钟级。为这 80% 的流量付跨 AZ 的复制费,本质上是用为核心链路设计的架构去承载非核心链路的流量

Diskless 给的就是一个「按 topic 分档」的能力。这才是它最值钱的地方——不是省了多少钱,而是让你能在同一个集群里同时做出两种不同的权衡

后续值得关注的方向

KIP-1150 的 Further Work 里列了几条,我认为其中两条会是接下来的重头戏:

  • Iceberg Format:允许 at-rest 的 topic 数据被大规模并行处理。这一步走通,Kafka topic 和数据湖表之间的边界就模糊了——你的消息队列直接就是一张 Iceberg 表,不需要 Connector、不需要 Flink 落湖。这可能比 diskless 本身影响更大。
  • Broker Roles:把 broker 按 produce / consume / coordination / compaction 分角色,允许异构集群。这是在数据面和元数据面分离之后的自然延伸——既然职责已经分开了,凭什么还要求每台机器规格一样?

另外还有 Topic Type Changing(classic ↔ diskless 互转)、Parallel Produce Handling(高延迟环境下提升单 producer 吞吐)、多 Region active-active。

时间线上,KIP-1150 已 Accepted,1163 和 1164 还在讨论。从 Kafka 的历史节奏看(KRaft 从提案到 production-ready 走了大约四年),Diskless 达到生产可用乐观估计也要到 2027 年

但架构方向已经定了。如果你现在正在做未来三年的流平台规划,「元数据与数据分离、块存储换对象存储」这个方向,值得现在就开始把它纳入设计约束——至少,别再往架构里塞更多依赖「leader 本地盘就是真相源」的假设了。


参考

  • KIP-1150: Diskless Topics(Accepted)
  • KIP-1163: Diskless Core(Under Discussion)
  • KIP-1164: Diskless Coordinator(Under Discussion)
  • KIP-405: Kafka Tiered Storage
  • KIP-630: Kafka Raft Snapshot
  • KIP-392: Allow consumers to fetch from closest replica
  • JIRA: KAFKA-19161

推荐文章

浏览器自动播放策略
2024-11-19 08:54:41 +0800 CST
在Vue3中实现代码分割和懒加载
2024-11-17 06:18:00 +0800 CST
一键配置本地yum源
2024-11-18 14:45:15 +0800 CST
Vue中的样式绑定是如何实现的?
2024-11-18 10:52:14 +0800 CST
最全面的 `history` 命令指南
2024-11-18 21:32:45 +0800 CST
JavaScript设计模式:发布订阅模式
2024-11-18 01:52:39 +0800 CST
程序员茄子在线接单