编程 Kafka Diskless 深度拆解:当 Broker 决定不再拥有磁盘——从 WAL Segment 到 SQLite 批次协调器的全链路架构手术

2026-08-11 05:50:06 +0800 CST views 4

Kafka Diskless 深度拆解:当 Broker 决定不再拥有磁盘——从 WAL Segment 到 SQLite 批次协调器的全链路架构手术

一句话概括:Apache Kafka 正在把「谁拥有数据」这个问题的答案,从 Broker 的本地磁盘,改写成对象存储;而把「谁决定顺序」这个问题,交给一个用 SQLite 做本地物化的新协调器。这不是加了个功能,这是把 Kafka 的存储层重新做了一遍。


一、先说结论:这次改的不是功能,是成本结构

如果你在 AWS 或者 GCP 上跑过一个稍微有点规模的 Kafka 集群,一定看过这样的账单:EC2 实例费用一栏还算克制,EBS 一栏有点肉疼,然后 Data Transfer - Regional 那一栏,数字大得让人怀疑人生。

我见过的典型比例是:跨可用区流量费能占到整个 Kafka 集群 TCO 的 40%~60%。你没看错,比机器本身还贵。

这钱是怎么烧掉的?很简单,因为经典 Kafka 的写入路径长这样:

Producer(AZ-a) ──► Leader(AZ-b)   ← 跨 AZ 第 1 次
                     │
                     ├──► Follower(AZ-c)  ← 跨 AZ 第 2 次
                     └──► Follower(AZ-a)  ← 跨 AZ 第 3 次

写 1 GiB,实际产生的跨 AZ 流量是「1(生产者到 Leader,如果没打中同区)+ 2(复制到另外两个副本)」= 最多 3 GiB。消费端如果没配 rack-aware fetch,再来一次。

按 KIP-1150 里明确列出的价格:

云厂商同 Region 跨 AZ 流量单价
AWS$0.02 / GiB
Google Cloud$0.01 / GiB
Azure不收费

注意 AWS 这 $0.02 是双向都收的口径下的常见理解,实际计费细节按官方文档为准。但即便按最保守的算法:

一个 100 MB/s 稳态入站的集群(这在互联网公司算中等规模),一个月入站数据量:

100 MB/s × 86400 s × 30 天 ≈ 259 TB/月

假设 RF=3,跨 AZ 复制系数按 2 算(Leader 到另外两个 AZ 的 Follower):

跨 AZ 复制流量 ≈ 259 TB × 2 = 518 TB
按 $0.02/GiB ≈ 518 × 1024 × 0.02 ≈ $10,600 / 月

这还没算消费端 fan-out。如果一份数据被 5 个消费组读,而且消费者没有 rack 亲和,跨 AZ 出口流量还能再翻几倍。

所以 Diskless Topics 要解决的核心问题只有一个:把这笔钱干掉。

它的手段也很直接——不再由 Kafka 自己做跨 AZ 复制,而是把数据一次性写进对象存储(S3 / GCS / Azure Blob),让云厂商的对象存储去负责跨 AZ 冗余。对象存储的 PUT/GET 在同 Region 内通常不收跨 AZ 流量费,冗余是内置的。


二、澄清一个致命误解:Diskless ≠ 没有磁盘

这是我看到最多的误读,KIP-1150 原文专门开了一节澄清。原话我翻译过来是:

Diskless 之于「无磁盘」,正如 Serverless 之于「无服务器」。

也就是说,磁盘还在,只是它不再是用户数据的持久化真相来源(source of truth)。Diskless topic 下,Broker 磁盘依然会被用来存这些东西:

  1. 常规 topic 的 KRaft 元数据 —— 集群元数据日志必须落盘,这个躲不掉;
  2. 批次元数据 —— 取决于批次协调器的实现方式,如果用 Kafka topic 实现,那它本身就要落盘;
  3. 正在上传到分层存储的 diskless 数据 —— 中转态;
  4. 给消费者做缓存的 diskless 数据 —— 这是性能优化,丢了不影响正确性。

第 4 点是关键:本地 segment 从「持久化机制」降级成了「缓存机制」。丢了就丢了,从对象存储重新拉一份就行。这个语义变化,是后面一系列架构变动的总开关。

你可以把 diskless topic 的本地保留期配到极短(比如几十秒),甚至在内存足够的集群上用内存盘顶上。


三、三代存储模型的演进路线图

理解 Diskless 之前,得先把 Kafka 存储模型的三代演进串起来。我画个对比表:

维度Classic TopicTiered Topic (KIP-405)Diskless Topic (KIP-1150)
活跃段存储位置本地块存储本地块存储对象存储
历史段存储位置本地块存储对象存储对象存储
复制方式Kafka 自己做(ISR)Kafka 自己做(活跃段)委托给对象存储
跨 AZ 复制流量有,且是大头有,仍是大头基本消除
Produce 必须打 Leader否,任意 Broker
偏移量分配者Leader BrokerLeader BrokerDiskless Coordinator
Produce P50 延迟个位数 ms ~ 几十 ms同左~500 ms
本地磁盘角色真相来源真相来源(活跃段)缓存
unclean leader election 风险消除

KIP-405 分层存储只解决了一半问题。 它把冷数据挪走了,省了存储费;但活跃段还在本地磁盘上,还得走 ISR 复制,跨 AZ 流量费一分没少。而对绝大多数 Kafka 集群来说,流量费才是大头,存储费反而是小头。

这就是为什么 WarpStream、AutoMQ、Confluent Freight 这些"Kafka 协议兼容替代品"能在市场上打出一片天——它们直接把活跃段也扔进了对象存储。

KIP-1150 的动机部分说得很直白(我意译):

这些替代方案正在获得市场成功,采用率在上升,说明市场对这个优化有普遍兴趣。上游 Kafka 应该把这个创新吸收进来。如果什么都不做,这将成为上游实现最大的功能缺口,会把高规模和云用户推向替代品,最终 Kafka 可能连协议的主导权都丢掉。

翻译成人话:再不做,协议就不是我们的了。 这是一次带着危机感的自我革命。


四、核心架构:数据与元数据的彻底解耦

4.1 Classic Produce 的六件事

在经典 topic 里,接收 Produce 请求的那个 Broker(必须是 Leader)要干六件事:

  1. 校验数据合法性
  2. 分配 offset 和 timestamp
  3. 把 offset/timestamp 注入到批次数据里
  4. 写入持久存储
  5. 等待复制完成
  6. 返回 Produce 响应

问题出在第 2、3 步:offset 是写进数据本身的。所以只有一个人能写,就是 Leader。这是 Kafka 单点写入模型的根。

4.2 Diskless 把职责拆开了

Diskless 的做法是:

  • 第 2 步(分配 offset/timestamp) → 交给 Diskless Coordinator
  • 第 3 步(注入 offset 到批次数据) → 交给 每个副本自己
  • 第 4 步(写持久存储) → 变成"上传对象到对象存储"
  • 第 5 步(等待复制) → 没了,对象存储自己保证

这个拆分带来的直接效果:任何 Broker 都能处理任何 diskless 分区的 Produce 请求。

producer 不用再关心谁是 Leader,直接找同 AZ 的那个 Broker写就行。跨 AZ 入站流量,消除。

4.3 顺序是怎么保证的?

这是最精妙的一刀。KIP-1163 原文有一句话我认为是整个设计的题眼:

我们把存储 I/O 和批次排序分离。持久化到对象存储的数据没有隐含顺序,顺序被委托给了批次坐标(batch coordinates)。

具体机制:

  • 每个 Broker 提交给 Coordinator 的坐标,只有局部顺序(相对于同一个 Broker 提交的其他坐标);
  • Diskless Coordinator 执行 append 动作,为每个分区里的批次坐标分配全局顺序

也就是说,对象存储里躺着的是一堆无序的字节块,真正的"日志顺序"活在 Coordinator 的元数据里。Broker 可以并行上传多个对象,但必须按顺序向 Coordinator 提交

这是典型的「数据平面无序、控制平面定序」设计。你在 Delta Lake / Iceberg 的事务日志里能看到同样的思路,只不过 Kafka 把它做到了毫秒级延迟的流式场景。


五、WAL Segment:为什么要把多个分区混装进一个对象

5.1 对象存储的成本模型逼出来的设计

经典 segment 有个铁律:一个 segment 只装一个分区的数据

但 WAL Segment 打破了这条规则——一个 WAL Segment 对象里混装了来自多个 topic、多个分区的批次

为什么要这么反直觉?KIP-1163 给的理由是"保持对象存储写操作成本合理"。展开说说这背后的账:

S3 的 PUT 请求单价大约 $0.005 / 1000 次(各 Region 略有差异)。假设一个集群有 10,000 个分区,如果每个分区单独一个对象,按每 250ms 刷一次:

每秒 PUT 次数 = 10000 分区 / 0.25 秒 = 40,000 次/秒
每月 PUT 次数 = 40000 × 86400 × 30 ≈ 1.04 × 10^11 次
每月费用 = 1.04e11 / 1000 × $0.005 ≈ $518,000

五十万美元一个月,只是为了发 PUT 请求。这个方案当场死亡。

改成混装之后:

每秒 PUT 次数 = Broker 数量 / 0.25 秒
30 个 Broker → 120 次/秒
每月 ≈ 120 × 86400 × 30 = 3.1 × 10^8 次
每月费用 ≈ 3.1e8 / 1000 × $0.005 ≈ $1,555

从 $518,000 降到 $1,555,差了 333 倍。

所以 KIP 原文那句"prevent the number of topic-partitions from influencing object read and write behavior"(防止分区数影响对象读写行为)不是随口说说,这是整个方案能不能成立的生死线。

这也顺带解决了 Kafka 一个老大难问题:分区数不再直接放大 IO 成本。 经典 Kafka 里,分区数暴涨会导致大量小文件、fsync 风暴、页缓存失效。Diskless 下,分区数对对象存储的压力几乎是常数级的。

5.2 WAL Segment 的物理格式

┌──────────────────────────────────────────────────────┐
│  Byte 0: version header (固定为 0)                    │
├──────────────────────────────────────────────────────┤
│  topic-A / partition-0 的批次区(连续)                 │
│    ├─ batch                                          │
│    ├─ batch                                          │
│    └─ batch                                          │
├──────────────────────────────────────────────────────┤
│  topic-A / partition-3 的批次区(连续)                 │
│    ├─ batch                                          │
│    └─ batch                                          │
├──────────────────────────────────────────────────────┤
│  topic-B / partition-7 的批次区(连续)                 │
│    └─ batch                                          │
└──────────────────────────────────────────────────────┘

几个关键设计点:

  1. 1 字节 header,当前固定为 0。KIP 明确说"解释后续批次数据不需要这个 header 里的信息"——这意味着老版本读取器可以跳过 1 字节直接解析,向前兼容留了口子。
  2. 按分区分组连续存放。这是为了局部性:后续要把某个分区的数据搬到分层存储时,可以用一次 ranged GET 把连续区间捞出来,而不是拆成 N 个小 GET。
  3. 对象名用全局唯一 ID(如 UUID),上传时无需协调、无冲突。这是能做到"任意 Broker 无锁并行上传"的前提。
  4. 对象不可变。一旦上传,内容与名字终身绑定。这既符合对象存储不支持随机写的现实,也让缓存策略极其简单——永不失效
  5. 批次的 offset 和 timestamp 可能是未设置的。真正的值由 Coordinator 在 commit 时确定,副本在往本地 segment 追加时注入。

第 5 点值得多说一句:这意味着对象存储里的字节流和消费者最终看到的字节流是不一样的。副本要做一次"注入",把 Coordinator 分配的 offset 写回批次头。这是保证 diskless 数据和 classic 数据在消费者眼里长得一模一样的关键——消费者不需要任何改动。


六、Produce 路径的延迟预算:500ms 是怎么来的

6.1 七步流程

① Producer 发 Produce 请求到【任意】Broker
        │
② Broker 把请求塞进本地 buffer,累积
        │  (超过大小或时间阈值才走下一步)
        │
③ Broker 生成 shared log segment + 所有缓冲批次的 batch coordinates
        │
④ 上传 WAL Segment 到对象存储(持久化完成)
        │
⑤ Broker 向 Diskless Coordinator 提交 batch coordinates
        │
⑥ Coordinator 分配 offset、持久化坐标、返回
        │
⑦ Broker 给这个对象关联的所有 Produce 请求返回响应

6.2 延迟拆解(KIP-1163 给的官方数字)

阶段P50P99可配置性
② 缓冲上限 250ms 或 4MiB可配置
④ 上传对象存储~100ms~200-400ms取决于云厂商
⑤⑥ 坐标提交~10ms~20-50ms取决于批次数量
端到端目标~500ms~1-2s

对比一下:经典 Kafka 在 acks=all 下,同 Region 的 Produce P99 通常在 10-50ms 量级。

Diskless 的延迟是经典模式的 20-100 倍。

这个数字必须摆在最前面说清楚,因为它直接决定了选型边界。任何跟你说"Diskless 是 Kafka 的未来,全都该换"的人,要么没算过延迟,要么没跑过线上。

6.3 三角权衡:你只能选两个

        低延迟
         /\
        /  \
       /    \
      /      \
     /        \
低成本 ────── 高吞吐
  • Classic Topic:低延迟 + 高吞吐,成本高
  • Diskless Topic:低成本 + 高吞吐,延迟高
  • 低延迟 + 低成本:物理上做不到(对象存储的往返延迟是硬约束)

KIP-1150 的高明之处在于:它把这个选择权交给了每个 topic,而不是每个集群。同一个集群里可以并存 classic、tiered、diskless 三种 topic,你按业务需求逐个 topic 决定。

这个"per-topic 而非 per-cluster"的决策,KIP 在 Rejected Alternatives 里专门解释过:集群是一个管理边界(统一的命名空间、权限体系、物理部署),同一个边界内的不同用户有不同的性能诉求,应该像 retention、segment rolling 一样可以逐 topic 配置。

我完全同意。这才是工程上正确的抽象层级。


七、Diskless Coordinator:用 SQLite 撑起 GB 级元数据

KIP-1164 是我认为整套设计里技术含量最高的部分。

7.1 它管什么

Diskless Coordinator(下称 DC)负责 diskless topic 的所有元数据:批次、WAL 文件、生产者状态、事务状态。它对外暴露 7 个 API(属于 Broker API,需要 CLUSTER_ACTION 权限,客户端不直接调用):

API作用一致性要求
DisklessCreatePartitions建分区/扩分区
DisklessDeleteTopics删 topic 及分区
DisklessCommitFile提交 WAL 文件 + 分配 offset(核心)
DisklessDeleteRecords删除分区尾部记录
DisklessListOffsets查 earliest/latest/按时间戳严格一致(必须反映所有前序变更)
DisklessFindBatches从指定 offset 查批次弱一致(可读陈旧状态、可由 Follower 服务)
DisklessDescribeFile查某 WAL 文件是否还有活批次

注意 DisklessListOffsetsDisklessFindBatches 的一致性差异——这是个很务实的设计。查 high watermark 必须准确(否则消费者会读到不该读的数据),但"从 offset 1000 开始给我找批次"这种查询,读到稍旧的视图完全没问题,大不了少返回几条。这个区分让读路径可以由 Follower 分担,直接缓解了 Coordinator 的热点问题。

7.2 分片:__diskless_metadata

DC 用一个新的内部 topic __diskless_metadata 做主存储,多分区。每个分区的 Leader 就是一个独立的 Coordinator。

  • 用户分区在创建时被分配到某个 DC,映射关系存在分区元数据里;
  • 各个 DC 在放置、操作上完全独立;
  • 可以给 __diskless_metadata 加分区来扩容,但只有新建的用户分区能用上新的元数据分区(这是个重要限制,后面踩坑清单会提);
  • 通过控制 __diskless_metadata 的分区放置,可以把 DC 集中到一部分专用 Broker 上——这实际上打开了"角色分离部署"的大门。

7.3 为什么是 SQLite

这是全文我最想聊的一个设计决策。

Kafka 现有的协调器(Group Coordinator、Transaction Coordinator)都是纯内存状态机 + Kafka topic 做日志。DC 为什么不能照抄?

KIP-1164 的原话(意译):

与 Kafka 里其他协调器不同,DC 的状态预期会很大(可达数百 MB 甚至数 GB),全放内存不现实。我们提议用 SQLite 在本地物化元数据日志。

想想这个规模是怎么来的:假设集群 10,000 个分区,每个分区活跃期内有 10,000 个批次坐标,每个坐标要记录 (partition, base_offset, wal_object_id, byte_offset, byte_length, timestamp, producer_id, epoch, sequence),保守按 100 字节算:

10000 × 10000 × 100 B = 10 GB

内存放不下。而且这还是每个 DC 副本都要有一份(Follower 也要物化,才能快速接管)。

SQLite 的价值在这里非常清晰:

  1. 结构化查询和索引开箱即用DisklessFindBatches 本质是个范围查询 WHERE partition = ? AND base_offset >= ? ORDER BY base_offset LIMIT ?,用 SQLite 的 B-tree 索引一行 SQL 搞定。自己实现一套磁盘友好的索引结构?那是又一个半年的工程量。
  2. ACID 事务。推测性执行需要"试着改、然后回滚",SQLite 的事务天然支持。
  3. 快照有现成方案。SQLite 的 Backup API 或只读事务能做点时间一致的异步快照。
  4. 它只是缓存,不是真相。KIP 明确说"SQLite DB 是本地元数据缓存,不是 source of truth,需要时可以直接删掉,从日志和快照重建"。

第 4 点是解除心理负担的关键:你不需要相信 SQLite 的持久性,你只需要相信它的查询能力。 真相永远在 __diskless_metadata 这个 Kafka topic 里。

这个"用嵌入式数据库物化复制日志"的模式,我认为是分布式系统里被严重低估的一招。etcd 用 bbolt,TiKV 用 RocksDB,现在 Kafka 用 SQLite。共同点是:Raft/日志负责一致性,本地引擎负责查询效率,两者职责清晰不越界。

7.4 推测性执行:流水线与恢复速度的矛盾调和

DC 处理一个变更操作的标准流程:

① 对当前状态做检查(能不能执行、全部还是部分)
② 生成一条或多条元数据记录,追加到本地日志
③ 等待记录复制到 Follower(等 high watermark 越过它)
④ 把已复制的记录真正应用到本地状态
⑤ 回复客户端

问题来了:③ 要等网络往返,如果串行等,吞吐直接崩。要做流水线,就得在上一个操作还没复制完的时候,开始处理下一个操作。但下一个操作的检查需要看到上一个操作的效果——而上一个操作还没应用到本地状态。

两个需求直接打架:

  • 要流水线 → 本地状态必须包含未复制的待定操作
  • 要快速恢复 → 本地状态只能包含已复制的操作(否则崩溃后可能要从头重建)

KIP-1164 的解法叫推测性选择性应用(speculatively selectively applies)

协调器把当前的待定操作推测性地应用到当前已提交的本地状态上,然后针对得到的推测状态做必要的检查。推测性应用可以完全在内存里做,也可以在一个注定要回滚的 SQLite 事务里做。

"注定要回滚的 SQLite 事务"(a to-be-rejected SQLite transaction)—— 这个说法我第一次读的时候笑出声,但仔细想,这是极其优雅的工程手段:

BEGIN;
  -- 把 pending 操作 1..N 应用上去
  INSERT INTO batches ...;
  UPDATE partitions SET high_watermark = ...;
  -- 针对推测状态做检查
  SELECT ... FROM batches WHERE producer_id = ? AND sequence = ?;
ROLLBACK;   -- 永远回滚,只借用它的隔离性做"假设推演"

用事务的隔离性做「if-then 推演」,检查完就丢弃。已提交状态永远干净,恢复时不需要回退任何东西。 两个矛盾需求同时满足。

7.5 日志膨胀怎么控

元数据日志最大的贡献者是批次和 WAL 文件的记录。DC 靠两件事控制规模:

  1. 与分层存储协同。周期性地把 diskless topic 的批次合并成标准 Kafka segment 卸载到 tiered storage。这意味着即使 topic 配置无限保留,单条批次元数据在 DC 里的生命周期也是有限的——由 segment.mssegment.bytes 决定。这个设计非常聪明:把无界问题转化成了有界问题。
  2. 快照 + 日志裁剪。对标 KRaft 的 KIP-630:Leader 异步做快照,Follower 拉快照,快照覆盖到的 offset 之前的日志可以裁掉。

八、副本语义的重新定义:Leader 也可能落后

这一节是我认为最容易踩坑的地方,因为它颠覆了 Kafka 十几年来的心智模型

8.1 ISR 的含义变了

经典 Kafka:in-sync = 与 Leader 同步。

Diskless:in-sync = 与 Diskless Coordinator 同步。

因为真相来源变成了 DC,所以:

Leader 从"同步"的角度看就是另一个普通副本,它也可能是 out-of-sync 的。

这句话的杀伤力很大,KIP-1163 列出了两个直接后果:

  • Leader 如果落后了,它没法高效履行自己的职责(比如卸载 segment 到分层存储);
  • 不建议从落后的 Leader 读(跟从任何落后副本读一样)。

8.2 为什么还要保留 Leader

既然不做复制了,为什么还要选举 Leader?KIP 给了明确答案,Leader 还要干三件事:

  1. 管理 ISR 状态
  2. 上传到分层存储
  3. 处理 share fetch(KIP-932 队列语义)

8.3 用空 FetchResponse 维持 ISR

这个细节很有意思:

尽管没有 Broker 间复制,副本仍然会向 Leader 发 FetchRequest。Leader 会返回空的(不含记录的) FetchResponse。这是 Leader 追踪 ISR 的机制。

也就是说,Fetch 协议被降级成了纯粹的心跳。副本真正的数据是从对象存储下载的,Fetch 只是为了让 Leader 知道"我还活着、我追到哪了"。

从工程角度看,这是典型的"复用现有机制降低改动面"——不用新造一套 ISR 心跳协议,把老协议掏空当心跳用。有点脏,但很有效。

8.4 unclean leader election 消失了

这是个被低估的收益。

经典 Kafka 里,unclean leader election(从落后副本里选 Leader)会导致数据丢失,这是运维最怕的场景之一。Diskless 下:

任何 Broker 都可以通过联系 Diskless Coordinator 来构建任意 diskless 分区的副本,降低其他 Broker 的负载,并消除 unclean leader election

因为数据在对象存储里,任何 Broker 随时可以从对象存储把状态补齐。"落后"只是暂时的缓存缺失,不是数据丢失。这个性质的改变,让 unclean.leader.election.enable 这个纠结了无数运维的配置项,在 diskless topic 上失去了意义。

8.5 副本重建策略

新副本上线时要从哪开始填本地日志?KIP 列了三种候选策略:

  1. 从最早的非 tiered offset 开始 —— 最完整,但可能很慢
  2. 从最新可用 offset 开始 —— 最快,但没有历史缓存
  3. 从当前 log end 往前推固定量(比如 100 MiB)或最早可用 offset,取较晚者 —— 折中

我倾向策略 3。策略 1 在大分区上会造成上线风暴(新 Broker 一上线就疯狂拉对象存储),策略 2 会让新副本在一段时间内对所有历史读都 cache miss。策略 3 有个"预热窗口"的概念,比较符合真实的消费者分布——绝大多数消费者都在读日志尾部。


九、WAL 文件的生命周期:一个链式移交的所有权协议

9.1 问题:一个文件,多个主人

WAL Segment 混装了多个分区的数据,而不同分区可能归属不同的 DC。那这个文件谁负责删?

如果 DC-1 说"我这边的批次都过期了,删文件吧",但 DC-5 那边还有活批次,删了就是数据丢失。

9.2 KIP-1164 的解法:随机排列的所有者链

① Broker 准备提交文件时,收集所有要提交的 DC 的 ID,组成 owners 列表
② 把 owners 列表【随机排列】
③ owners 列表作为字段放进 DisklessCommitFile 请求 —— 现在每个 DC 都知道有哪些 DC 声称拥有这个文件
④ owners[0] 成为当前所有者
⑤ 当 owners[0] 里这个文件的最后一个批次被删除时,把文件【移交】给 owners[1]
⑥ 一直传递下去,最后一个所有者删完时,物理删除文件

为什么要随机排列? KIP 没展开,但我认为答案很清楚:避免热点。如果不随机,比如总是按 DC ID 升序,那么 DC-0 会成为绝大多数文件的首任所有者,承担所有的初始所有权管理开销和后台移交工作。随机排列把这个负担均摊到所有 DC 上。

这是分布式系统里"用随机化打散热点"的经典手法,代价是一行代码,收益是消除一整类倾斜问题。

9.3 移交和删除必须异步

KIP 强调移交和删除不能干扰 DC 的正常活动,必须异步:

  • 文件状态变更只改本地状态不写元数据日志(因为这个信息已经隐含在日志里了,Follower 自己能推导出来);
  • 启动后台 worker 执行移交或删除;
  • 目标 Broker 或对象存储不可用时,后台 worker 无限重试
  • DC leadership 变更时,当前 Broker 上的 worker 停止,在新 Leader 上重启(因为新 Leader 读的是同一份日志,知道文件状态)。

"不写元数据日志"这一点值得学习——能从现有日志推导出来的状态,就不要再写一遍日志。这是控制日志膨胀的基本功。

9.4 孤儿文件回收

Broker 上传了文件但没提交成功(比如上传完就崩了),就产生孤儿文件。清理算法:

① 每次扫描前,向每个 DC 询问"你手上最老的、还有批次的文件的时间戳"
② 取所有 DC 里最老的那个时间戳,加上 grace period,作为阈值
③ 扫描对象存储,找比阈值更老的文件
④ 【额外安全措施】再向每个 DC 确认这些具体文件是否已知
⑤ 任何 DC 都不认识的文件 → 物理删除

第 ④ 步的"额外安全措施"是关键。第 ②③ 步的时间戳判断是启发式的,可能因为时钟漂移、DC 短暂不可用等原因误判。第 ④ 步做一次精确核对,把误删概率压到极低。

扫描频率必须可配置,默认应该很低——孤儿文件是罕见事件,为它频繁 LIST 对象存储(LIST 请求比 GET 贵得多)不划算。


十、Produce Gateway:把 N 次提交压到 1 次

10.1 提交放大问题

"任何 Broker 都能处理任何分区的 Produce"这个自由度,带来了一个副作用:

一个 WAL 文件里如果混装了归属 12 个不同 DC 的分区数据,Broker 就要发 12 次 DisklessCommitFile 网络调用。

最坏情况下,出站提交调用数 = n_dcs(DC 数量)。这会导致:

  • 部分失败概率上升(12 个调用里挂一个就麻烦)
  • 尾延迟恶化(要等最慢的那个)

部分缓解:一个 Broker 可以同时托管多个 DC,逻辑提交可以合并成更少的物理请求。但最坏情况上界还是 n_brokers

10.2 解法:PreferredProduceBrokers

KIP-1163 扩展了 Metadata 请求/响应,为新客户端增加 PreferredProduceBrokers 字段,值由 Broker 侧动态计算。

机制:为每个 DC,在每个 rack 里选一个或几个 Broker 作为"produce gateway"。生产者要往某个 DC 管理的 diskless 分区写时,会被引导到本 rack 里对应的 gateway Broker,而不是随便找一个。

KIP 里给了两个场景:

场景 1:Broker 数少于 DC 数

Racks: 3, Brokers: 3, DCs: 12
每个 DC 在每个 rack 选 1 个 gateway
→ 每个 Broker 是 12 × 3 / 3 = 12 个 DC 的 gateway
→ 但因为 Broker 比 DC 少,每个 Broker 每个 WAL 文件最多做 n_brokers 次出站调用

场景 2:Broker 数充足

只要 Broker 够多,可以做到每个 Broker 在其 rack 内只做 1 个 DC 的 gateway
→ 每个 WAL 文件只需 1 次提交调用

这个设计的本质是:用客户端路由的确定性,换取服务端提交的局部性。 跟 Kafka 原本的"生产者按分区路由到 Leader"是同构的思路,只是路由目标从"分区 Leader"变成了"DC gateway"。

注意:这需要新版客户端。 老客户端不认识 PreferredProduceBrokers,会退化成随机打,提交放大问题依然存在。这是升级路径上必须考虑的因素。


十一、代码实战

下面这部分是我基于 KIP 设计推演出来的实践方案。注意:Diskless 在写作时仍处于 KIP-1163/1164 Under Discussion 阶段,具体配置项名称以最终实现为准,这里重点是让你理解调参逻辑。

11.1 生产者侧:为高延迟环境调参

Diskless 的 P50 ~500ms 意味着经典的生产者配置全部失效。核心矛盾:单个生产者的吞吐 = in-flight 请求数 × 每请求大小 / RTT。RTT 从 10ms 涨到 500ms,吞吐直接掉 50 倍——除非你把 in-flight 提上去。

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;

public class DisklessProducerConfig {

    public static Properties build(String bootstrap, String rack) {
        Properties p = new Properties();
        p.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap);
        p.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        p.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);

        // ── 1. 幂等必开 ──────────────────────────────────────
        // KIP-1163 明确说:上传成功但提交失败会留下垃圾对象,
        // 生产者重试是安全的【前提是使用幂等 produce】。
        // 非幂等模式下重试 = 数据重复。这是硬性要求,不是建议。
        p.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
        p.put(ProducerConfig.ACKS_CONFIG, "all");

        // ── 2. in-flight 深度:对抗高 RTT 的唯一手段 ────────────
        // 经典配置常见 5。Diskless 下 RTT 涨 50 倍,
        // in-flight 不提上来,单生产者吞吐会被 RTT 直接锁死。
        // 幂等模式下 Kafka 保证乱序重试也能维持分区内顺序。
        p.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);

        // ── 3. linger:与 Broker 侧缓冲对齐 ──────────────────
        // Broker 侧会缓冲最多 250ms。客户端再 linger 100ms,
        // 换来更大的批次 → 更少的坐标 → 更快的 commit。
        // 反正端到端已经 500ms 了,这 100ms 是"免费"的。
        p.put(ProducerConfig.LINGER_MS_CONFIG, 100);
        p.put(ProducerConfig.BATCH_SIZE_CONFIG, 1024 * 1024);   // 1 MiB

        // ── 4. 超时:必须放大,否则正常延迟被误判为失败 ──────────
        // P99 目标 1-2s,request.timeout 给到 30s 留足余量。
        // delivery.timeout 覆盖整个重试窗口。
        p.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, 30_000);
        p.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 300_000);
        p.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);
        p.put(ProducerConfig.RETRY_BACKOFF_MS_CONFIG, 200);

        // ── 5. buffer.memory:高延迟下必须放大 ───────────────
        // 未确认数据在内存里滞留的时间变长了 50 倍,
        // 缓冲区按原来配就会频繁触发阻塞。
        // 粗算:目标吞吐 × 端到端延迟 × 安全系数
        //      50 MB/s × 2s × 2 = 200 MB
        p.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 256L * 1024 * 1024);
        p.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, 60_000);

        // ── 6. rack 亲和:省钱的关键 ─────────────────────────
        // 让客户端能被引导到同 AZ 的 produce gateway
        p.put(ProducerConfig.CLIENT_ID_CONFIG, "diskless-producer-" + rack);
        p.put("client.rack", rack);

        // ── 7. 压缩:CPU 换对象存储成本 ──────────────────────
        // 对象存储按 GB 计费 + PUT 按次计费,
        // zstd 压缩率通常比 lz4 高 20-30%,直接省存储费。
        // 高延迟场景下压缩那点 CPU 时间完全被 IO 淹没。
        p.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd");
        p.put("compression.zstd.level", 3);

        return p;
    }
}

关键点解释:

  • enable.idempotence=true 是硬性要求。KIP 原文说得很清楚,"生产者重试是安全的,前提是它们使用幂等 produce"。因为上传成功、提交失败的场景下,客户端会收到错误并重试,但对象已经在存储里了。没有幂等保护,重试就是重复数据。
  • buffer.memory 的计算方式目标吞吐 × 端到端延迟 × 安全系数。这是个小定律(Little's Law)的应用——系统内的在途数据量 = 到达率 × 停留时间。延迟涨 50 倍,在途数据量就涨 50 倍。

11.2 一个延迟预算探针

上线前你必须知道自己的延迟分布长什么样。这个探针程序做的是分段计时,把延迟归因到具体环节:

import org.apache.kafka.clients.producer.*;
import org.HdrHistogram.Histogram;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicLong;

/**
 * Diskless 延迟预算探针。
 * 输出 P50/P90/P99/P999,用于判断是否符合 KIP 预期
 * (P50 ~500ms, P99 ~1-2s)。偏离过大说明配置或环境有问题。
 */
public class DisklessLatencyProbe {

    // HdrHistogram:最大 60s,3 位有效数字
    private final Histogram e2e = new Histogram(60_000_000_000L, 3);
    private final AtomicLong errors = new AtomicLong();
    private final AtomicLong sent   = new AtomicLong();

    public void run(String topic, Properties props,
                    int durationSec, int ratePerSec) throws Exception {

        try (Producer<String, byte[]> producer = new KafkaProducer<>(props)) {
            ScheduledExecutorService sched =
                Executors.newScheduledThreadPool(4);

            byte[] payload = new byte[1024];   // 1 KiB 消息
            long intervalNs = 1_000_000_000L / ratePerSec;

            CountDownLatch done = new CountDownLatch(1);
            long endAt = System.nanoTime() + durationSec * 1_000_000_000L;

            sched.scheduleAtFixedRate(() -> {
                if (System.nanoTime() > endAt) { done.countDown(); return; }

                final long t0 = System.nanoTime();
                sent.incrementAndGet();

                producer.send(
                    new ProducerRecord<>(topic, null, payload),
                    (meta, ex) -> {
                        if (ex != null) {
                            errors.incrementAndGet();
                            return;
                        }
                        long elapsed = System.nanoTime() - t0;
                        synchronized (e2e) { e2e.recordValue(elapsed); }
                    });

            }, 0, intervalNs, TimeUnit.NANOSECONDS);

            done.await(durationSec + 30, TimeUnit.SECONDS);
            sched.shutdownNow();
            producer.flush();
        }

        report();
    }

    private void report() {
        System.out.println("─────── Diskless 延迟预算报告 ───────");
        System.out.printf("发送总数   : %d%n", sent.get());
        System.out.printf("失败数     : %d (%.4f%%)%n",
                errors.get(), 100.0 * errors.get() / Math.max(1, sent.get()));
        System.out.printf("P50        : %8.1f ms   [KIP 目标 ~500]%n", ms(50.0));
        System.out.printf("P90        : %8.1f ms%n", ms(90.0));
        System.out.printf("P99        : %8.1f ms   [KIP 目标 1000-2000]%n", ms(99.0));
        System.out.printf("P99.9      : %8.1f ms%n", ms(99.9));
        System.out.printf("Max        : %8.1f ms%n", e2e.getMaxValue() / 1e6);

        // 归因提示
        double p50 = ms(50.0);
        if (p50 < 300) {
            System.out.println("⚠  P50 明显低于预期 —— 确认这个 topic 真的是 diskless,"
                             + "别是建成 classic 了");
        } else if (p50 > 900) {
            System.out.println("⚠  P50 显著偏高 —— 排查顺序:"
                             + "①对象存储同 Region 吗 ②DC 是不是热点 "
                             + "③客户端 linger 是不是配太大");
        } else {
            System.out.println("✓  P50 在预期区间内");
        }

        double ratio = ms(99.0) / Math.max(1.0, p50);
        if (ratio > 5.0) {
            System.out.printf("⚠  P99/P50 = %.1f,尾延迟发散。"
                    + "典型原因:单个 WAL 文件跨了太多 DC,"
                    + "提交要等最慢的那个。检查 PreferredProduceBrokers 是否生效。%n", ratio);
        }
    }

    private double ms(double pct) {
        synchronized (e2e) {
            return e2e.getValueAtPercentile(pct) / 1e6;
        }
    }

    public static void main(String[] args) throws Exception {
        Properties p = DisklessProducerConfig.build(
            System.getenv().getOrDefault("BOOTSTRAP", "localhost:9092"),
            System.getenv().getOrDefault("RACK", "az-a"));
        new DisklessLatencyProbe().run("probe-diskless", p, 120, 500);
    }
}

P99/P50 > 5 这个判据是我加的经验规则。按 KIP 给的数字,正常情况下这个比值应该在 2-4 之间(500ms → 1-2s)。如果发散得厉害,最可能的原因就是提交放大——一个 WAL 文件跨了太多 DC,要等最慢的那个返回。

11.3 成本对比脚本

选型决策要用数字说话。这个脚本把 classic 和 diskless 的月成本算出来:

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
Kafka Classic vs Diskless 月度成本估算器。
免责声明:单价随云厂商/Region 变动,这里只做数量级对比。
"""
from dataclasses import dataclass


@dataclass
class Pricing:
    cross_az_per_gib: float = 0.02      # AWS 同 Region 跨 AZ
    ebs_gp3_per_gb_month: float = 0.08
    s3_standard_per_gb_month: float = 0.023
    s3_put_per_1k: float = 0.005
    s3_get_per_1k: float = 0.0004


@dataclass
class Workload:
    ingress_mb_per_sec: float           # 稳态入站
    retention_hours: float              # 保留时长
    replication_factor: int = 3
    consumer_groups: int = 3            # fan-out 倍数
    num_partitions: int = 1000
    num_brokers: int = 12
    az_count: int = 3
    # 消费者中有多大比例做到了 rack 亲和(0~1)
    consumer_rack_affinity: float = 0.0
    # diskless 侧参数
    wal_flush_interval_sec: float = 0.25
    compression_ratio: float = 0.35     # zstd 后剩余比例


SEC_PER_MONTH = 86400 * 30
GIB = 1024


def monthly_ingress_gib(w: Workload) -> float:
    return w.ingress_mb_per_sec * SEC_PER_MONTH / GIB


def retained_gib(w: Workload) -> float:
    return w.ingress_mb_per_sec * w.retention_hours * 3600 / GIB


def classic_cost(w: Workload, p: Pricing) -> dict:
    ing = monthly_ingress_gib(w)

    # 复制:Leader 发给 RF-1 个副本,假设都跨 AZ
    repl_gib = ing * (w.replication_factor - 1)

    # 生产者入站:假设 1/az_count 概率命中同 AZ 的 Leader
    prod_cross = ing * (1 - 1.0 / w.az_count)

    # 消费出站
    egress = ing * w.consumer_groups * (1 - w.consumer_rack_affinity) \
             * (1 - 1.0 / w.az_count)

    net_cost = (repl_gib + prod_cross + egress) * p.cross_az_per_gib

    # 存储:每个副本一份,EBS 还要留 buffer(按 1.4 倍)
    storage_gb = retained_gib(w) * w.replication_factor * 1.4
    storage_cost = storage_gb * p.ebs_gp3_per_gb_month

    return {
        "跨AZ-复制": repl_gib * p.cross_az_per_gib,
        "跨AZ-生产": prod_cross * p.cross_az_per_gib,
        "跨AZ-消费": egress * p.cross_az_per_gib,
        "块存储": storage_cost,
        "合计": net_cost + storage_cost,
    }


def diskless_cost(w: Workload, p: Pricing) -> dict:
    ing = monthly_ingress_gib(w)
    stored = ing * w.compression_ratio

    # 跨 AZ:生产者被引导到同 AZ 的 gateway → ~0
    # Broker → S3 同 Region → 不计跨 AZ
    # 消费:热尾读命中本地缓存,冷读走 S3(同 Region)
    net_cost = 0.0

    # S3 存储(只存一份,冗余由 S3 内部负责)
    storage_gb = retained_gib(w) * w.compression_ratio
    storage_cost = storage_gb * p.s3_standard_per_gb_month

    # PUT:每个 Broker 每 wal_flush_interval 一次
    puts = w.num_brokers / w.wal_flush_interval_sec * SEC_PER_MONTH
    put_cost = puts / 1000 * p.s3_put_per_1k

    # GET:粗估冷读比例 20%,按 4 MiB 一次 ranged get
    cold_gib = ing * w.consumer_groups * 0.2
    gets = cold_gib * GIB / 4
    get_cost = gets / 1000 * p.s3_get_per_1k

    # 元数据日志(__diskless_metadata)仍走经典复制,但量很小
    # 粗估为入站量的 0.5%
    meta_cost = ing * 0.005 * (w.replication_factor - 1) * p.cross_az_per_gib

    return {
        "跨AZ-复制": 0.0,
        "跨AZ-生产": 0.0,
        "跨AZ-消费": 0.0,
        "对象存储": storage_cost,
        "PUT请求": put_cost,
        "GET请求": get_cost,
        "元数据日志": meta_cost,
        "合计": storage_cost + put_cost + get_cost + meta_cost,
    }


def compare(w: Workload, p: Pricing = Pricing()):
    c, d = classic_cost(w, p), diskless_cost(w, p)
    print(f"\n{'='*58}")
    print(f"负载:{w.ingress_mb_per_sec} MB/s 入站, "
          f"保留 {w.retention_hours}h, RF={w.replication_factor}, "
          f"{w.consumer_groups} 个消费组, {w.num_partitions} 分区")
    print(f"月入站量:{monthly_ingress_gib(w)/1024:.1f} TiB")
    print('='*58)

    print(f"\n【Classic】")
    for k, v in c.items():
        mark = "  ← 大头" if k.startswith("跨AZ") and v > c["合计"] * 0.25 else ""
        print(f"  {k:<12} ${v:>12,.0f}{mark}")

    print(f"\n【Diskless】")
    for k, v in d.items():
        print(f"  {k:<12} ${v:>12,.0f}")

    saving = c["合计"] - d["合计"]
    pct = saving / c["合计"] * 100 if c["合计"] else 0
    print(f"\n>>> 月度节省 ${saving:,.0f} ({pct:.1f}%)")
    print(f">>> 年度节省 ${saving*12:,.0f}")
    print(f">>> 代价:Produce P50 从 ~10ms 涨到 ~500ms(50 倍)")


if __name__ == "__main__":
    # 场景 A:日志/埋点收集 —— diskless 的甜蜜区
    compare(Workload(ingress_mb_per_sec=100, retention_hours=72,
                     consumer_groups=3, num_partitions=2000))

    # 场景 B:小集群 —— 省的钱可能不值得延迟代价
    compare(Workload(ingress_mb_per_sec=5, retention_hours=24,
                     consumer_groups=2, num_partitions=100,
                     num_brokers=3))

    # 场景 C:已经做好 rack 亲和的集群 —— 收益打折
    compare(Workload(ingress_mb_per_sec=100, retention_hours=72,
                     consumer_groups=3, num_partitions=2000,
                     consumer_rack_affinity=0.9))

跑一下你会发现几个反直觉的结论:

  1. 小集群(<10 MB/s)迁 diskless 基本不值。省的几百美元换 50 倍延迟,不划算。
  2. 已经做好 rack-aware consumer 的集群,收益会打折——因为消费侧的跨 AZ 成本你已经优化过了。剩下的主要是复制侧的收益。
  3. 压缩率对 diskless 的影响远大于对 classic 的影响。因为 classic 的成本大头是流量(压缩后也省),diskless 的成本大头是存储(压缩后直接省)。zstd 值得开。

11.4 本地实验环境

想在本机验证 diskless 语义,用 MinIO 起个 S3 兼容存储:

# docker-compose.yml —— Diskless 实验环境
version: "3.9"

services:
  minio:
    image: minio/minio:latest
    command: server /data --console-address ":9001"
    environment:
      MINIO_ROOT_USER: minioadmin
      MINIO_ROOT_PASSWORD: minioadmin
    ports: ["9000:9000", "9001:9001"]
    volumes: ["minio-data:/data"]
    healthcheck:
      test: ["CMD", "mc", "ready", "local"]
      interval: 5s
      timeout: 3s
      retries: 20

  minio-init:
    image: minio/mc:latest
    depends_on:
      minio: { condition: service_healthy }
    entrypoint: >
      /bin/sh -c "
      mc alias set local http://minio:9000 minioadmin minioadmin;
      mc mb -p local/kafka-diskless;
      mc anonymous set none local/kafka-diskless;
      echo '✓ bucket kafka-diskless ready';
      "

  # 用 toxiproxy 注入延迟,模拟真实云环境的对象存储 RTT
  # 本机 MinIO 只有 1ms 延迟,测不出真实的延迟预算
  toxiproxy:
    image: ghcr.io/shopify/toxiproxy:latest
    ports: ["8474:8474", "9010:9010"]
    depends_on: [minio]

volumes:
  minio-data:

配套的延迟注入脚本——这一步很多人会跳过,但它是本地实验有没有意义的分水岭

#!/usr/bin/env bash
# inject-s3-latency.sh
# 本机 MinIO 的 RTT ≈ 1ms,云上 S3 的 P50 ≈ 100ms、P99 ≈ 400ms。
# 不注入延迟的本地测试,跑出来的数字毫无参考价值。
set -euo pipefail

PROXY_API="http://localhost:8474"

echo "→ 创建 S3 代理"
curl -s -X POST "$PROXY_API/proxies" -d '{
  "name": "s3",
  "listen": "0.0.0.0:9010",
  "upstream": "minio:9000",
  "enabled": true
}' >/dev/null

echo "→ 注入上行延迟:基线 100ms,抖动 ±150ms"
curl -s -X POST "$PROXY_API/proxies/s3/toxics" -d '{
  "name": "latency_up",
  "type": "latency",
  "stream": "upstream",
  "attributes": { "latency": 100, "jitter": 150 }
}' >/dev/null

echo "→ 注入下行延迟:基线 50ms,抖动 ±50ms"
curl -s -X POST "$PROXY_API/proxies/s3/toxics" -d '{
  "name": "latency_down",
  "type": "latency",
  "stream": "downstream",
  "attributes": { "latency": 50, "jitter": 50 }
}' >/dev/null

echo "✓ 现在把 Kafka 的 S3 endpoint 指向 localhost:9010"
echo ""
echo "混沌测试选项:"
echo "  # 模拟对象存储间歇性不可用(测试 orphan file 回收)"
echo "  curl -X POST $PROXY_API/proxies/s3 -d '{\"enabled\":false}'"
echo "  sleep 30"
echo "  curl -X POST $PROXY_API/proxies/s3 -d '{\"enabled\":true}'"
echo ""
echo "  # 模拟带宽受限(测试上传背压)"
echo "  curl -X POST $PROXY_API/proxies/s3/toxics -d '{"
echo "    \"name\":\"bandwidth\",\"type\":\"bandwidth\",\"stream\":\"upstream\","
echo "    \"attributes\":{\"rate\":10240}}'   # 10 MB/s"

为什么必须注入延迟:整套 Diskless 设计的性能特征完全由对象存储的 RTT 决定。本机 MinIO 的 1ms RTT 会让你得出"Diskless 延迟也就 30ms 嘛"的错误结论,然后上云被打脸。

11.5 消费者侧:Go 客户端与 rack 亲和

package main

import (
	"context"
	"fmt"
	"log"
	"os"
	"time"

	"github.com/twmb/franz-go/pkg/kgo"
)

// Diskless topic 的消费者配置。
// 核心目标:尽量命中【同 AZ 副本的本地缓存】,
// 避免走对象存储冷读(既慢又要付 GET 费用)。
func newDisklessConsumer(brokers []string, group, topic, rack string) (*kgo.Client, error) {
	opts := []kgo.Option{
		kgo.SeedBrokers(brokers...),
		kgo.ConsumerGroup(group),
		kgo.ConsumeTopics(topic),

		// ① rack 亲和:优先从同 AZ 副本读。
		//    Diskless 下这一条依然重要——虽然数据在对象存储,
		//    但副本本地缓存能省掉一次 S3 GET(省钱 + 省延迟)。
		kgo.Rack(rack),

		// ② 拉大 fetch 批量。
		//    Diskless 的批次天然更大(Broker 侧缓冲了 250ms/4MiB),
		//    小 fetch 会造成频繁往返。
		kgo.FetchMinBytes(1 << 20),          // 1 MiB
		kgo.FetchMaxBytes(50 << 20),         // 50 MiB
		kgo.FetchMaxWait(500 * time.Millisecond),

		// ③ 超时放大:冷读要走对象存储,可能很慢
		kgo.RequestTimeoutOverhead(30 * time.Second),
		kgo.RetryTimeout(2 * time.Minute),

		// ④ 手动提交:diskless 端到端延迟高,
		//    自动提交容易在 rebalance 时丢/重
		kgo.DisableAutoCommit(),
	}
	return kgo.NewClient(opts...)
}

// 带冷读检测的消费循环
func consume(cl *kgo.Client) {
	ctx := context.Background()

	// 冷读检测:如果单次 poll 耗时远超 FetchMaxWait,
	// 大概率是走了对象存储冷读
	const coldReadThreshold = 2 * time.Second

	var (
		total     int64
		coldReads int64
		lastLog   = time.Now()
	)

	for {
		t0 := time.Now()
		fetches := cl.PollRecords(ctx, 10000)
		elapsed := time.Since(t0)

		if errs := fetches.Errors(); len(errs) > 0 {
			for _, e := range errs {
				log.Printf("fetch error topic=%s partition=%d: %v",
					e.Topic, e.Partition, e.Err)
			}
			continue
		}

		n := fetches.NumRecords()
		if n == 0 {
			continue
		}

		if elapsed > coldReadThreshold {
			coldReads++
			log.Printf("⚠ 疑似冷读:poll 耗时 %v,取回 %d 条。"+
				"若频繁出现,考虑增大副本本地缓存保留期", elapsed, n)
		}

		fetches.EachRecord(func(r *kgo.Record) {
			total++
			_ = r // 你的业务处理
		})

		// 手动提交
		if err := cl.CommitUncommittedOffsets(ctx); err != nil {
			log.Printf("commit failed: %v", err)
		}

		if time.Since(lastLog) > 10*time.Second {
			ratio := float64(coldReads) / float64(total+1) * 100
			log.Printf("已消费 %d 条,冷读事件 %d 次 (%.3f%%)",
				total, coldReads, ratio)
			lastLog = time.Now()
		}
	}
}

func main() {
	rack := os.Getenv("AZ")
	if rack == "" {
		rack = "az-a"
	}
	cl, err := newDisklessConsumer(
		[]string{"kafka-1:9092", "kafka-2:9092", "kafka-3:9092"},
		"diskless-consumer-group", "events-diskless", rack)
	if err != nil {
		log.Fatal(err)
	}
	defer cl.Close()

	fmt.Printf("消费者启动,rack=%s\n", rack)
	consume(cl)
}

冷读检测那段是我加的实用逻辑。Diskless 下消费成本的关键变量是缓存命中率:命中本地缓存 = 免费且快,miss 走对象存储 = 付 GET 费用且慢几百毫秒。生产环境应该把这个指标接进监控。


十二、选型决策树:哪些 topic 该换,哪些不该

我把判断逻辑整理成一棵树:

这个 topic 的端到端延迟 SLA 是多少?
│
├─ < 100ms(撮合、风控、实时竞价、RPC-over-Kafka)
│     → ❌ 绝对不要 Diskless。500ms P50 直接违约。
│
├─ 100ms ~ 1s(用户可感知的准实时,如实时推荐特征)
│     → ⚠️ 谨慎。P99 可能到 1-2s,会击穿 SLA。
│        建议先做全链路压测,或只在非核心链路试点。
│
└─ > 5s(日志、埋点、审计、CDC 落湖、指标、数仓入湖)
      → ✅ Diskless 的甜蜜区
      │
      └─ 吞吐量有多大?
            ├─ < 10 MB/s → ⚠️ 省的钱可能不够折腾成本,先算账
            └─ > 50 MB/s → ✅ 强烈建议,收益显著

具体场景对照表:

场景建议理由
应用日志收集✅ Diskless延迟不敏感、量巨大、成本敏感
用户行为埋点✅ Diskless同上,且通常有 T+0 小时级下游
CDC 数据入湖✅ Diskless落湖本来就是分钟级,且能配合 Iceberg 后续优化
监控指标✅ Diskless采集周期本来就是 10-60s
审计日志✅ Diskless强调持久性而非延迟,对象存储持久性更高
微服务异步解耦⚠️ 看情况如果调用方在同步等结果,500ms 可能不可接受
订单状态流转⚠️ 谨慎用户会盯着页面刷新
Kafka Streams 多级拓扑❌ Classic延迟会逐级累加,3 级拓扑 = 1.5s+
支付/撮合❌ Classic无需讨论
请求-响应模式❌ Classic一来一回 1s,用户直接跑了

Kafka Streams 那条要特别强调:如果你有 source → transform → repartition → aggregate → sink 这样的多级拓扑,每一级 repartition 都是一次完整的 produce+consume 往返。5 级拓扑在 diskless 下端到端延迟能到 2.5-5 秒。这个坑我认为会是社区未来一年里最常见的事故来源。


十三、十六条踩坑清单

按我理解的严重程度排序:

1. 不开幂等就上 Diskless = 定期数据重复。
KIP 明确:上传成功、提交失败会留下垃圾对象并返回错误给生产者,重试是安全的前提是幂等 produce。非幂等下每次这种失败都可能产生重复。这不是"最佳实践",是硬性前提。

2. __diskless_metadata 扩分区不会惠及已有 topic。
KIP 白纸黑字:加分区只有新创建的用户分区能用上新的元数据分区/DC。这意味着初始分区数是个长期决策,配少了会成为长期瓶颈。上线前务必按未来 2-3 年的分区规模规划。DC 之间的用户分区迁移在 KIP-1164 里明确标注为 out of scope。

3. Leader 可能是 out-of-sync 的,别再假设"读 Leader 最新"。
这是最反直觉的一条。经典 Kafka 里"Leader 必然拥有最新数据"是公理,Diskless 下它不成立了。KIP 明确说"不建议从落后的 Leader 读"。所有依赖"Leader 一定最新"的运维脚本、监控告警、故障处理 SOP 都要重写。

4. 老客户端会导致提交放大。
PreferredProduceBrokers 是 Metadata 协议的扩展,只对新客户端生效。混合了老客户端的集群,最坏情况一个 WAL 文件要提交 n_brokers 次,尾延迟会很难看。客户端升级要排在 Diskless 启用之前

5. Kafka Streams 多级拓扑延迟叠加。
见上一节。上线前把拓扑图画出来,数一数有几次 repartition,乘以 500ms。

6. 本地测试不注延迟 = 白测。
MinIO 本机 1ms RTT,S3 真实 100-400ms。差两个数量级,所有结论都不成立。用 toxiproxy 或类似工具注入。

7. buffer.memory 不放大会频繁阻塞。
在途数据量按 Little's Law 涨了 50 倍。原来 32MB 够用的,现在要 256MB+。不改的话 send() 会频繁阻塞在 max.block.ms 上。

8. request.timeout.ms 不放大会误判失败。
默认 30s 其实还行,但很多团队为了快速失败调到了 3-5s。Diskless 下 P99 就是 1-2s,加上重试和抖动,5s 超时会大量误报。

9. 孤儿文件扫描频率配高了会烧钱。
LIST 请求比 GET 贵。KIP 建议"默认扫描频率应该相应地低"。别为了强迫症把它配成每分钟一次。

10. WAL 文件删除的 grace period 配短了会打断消费者冷读。
KIP 说"文件要删除时,应该允许一个宽限期,不要立刻删,让消费者完成可能正在进行的读"。配得太短,正在做冷读的消费者会拿到 404。

11. 混合 topic 类型的 Produce 请求会被拖慢。
KIP-1163 明确:一个 Produce 请求里如果同时包含 classic 和 diskless 分区,响应会被延迟到 diskless 分区提交完成。也就是说,一个 diskless 分区会把整个请求里的 classic 分区一起拖到 500ms。生产者应该把 classic 和 diskless 的写入分开到不同的 Producer 实例

12. 事务只支持 v2 及以上。
KIP-1163 明确说"我们提议只关注 v2 及更新的事务版本"。老版本事务客户端不能用。

13. 副本重建策略选错会造成上线风暴。
选了"从最早非 tiered offset 开始",新 Broker 一上线就疯狂拉对象存储,可能把带宽打满、GET 费用飙升。大集群建议用"往前推固定量"的折中策略。

14. DC 状态 GB 级,磁盘规划要跟上。
SQLite 库会长到 GB 级,而且每个 DC 副本都有一份。托管 DC 的 Broker 需要额外的本地磁盘空间和 IOPS 预算。别用最小规格实例托管 DC。

15. 压缩率直接影响成本,别用 none。
Diskless 的成本大头从流量变成了存储,压缩的收益被放大了。zstd level 3 是个很好的起点。

16. 别指望 Azure 用户有同样的收益。
Azure 不收跨 AZ 流量费。Azure 上的 Diskless 收益主要来自运维简化(更容易扩缩容、对象存储持久性更高),而不是成本。要重新算账。


十四、上游收编:一场关于协议主导权的阳谋

最后聊点技术之外的东西。

KIP-1150 的 "Do Nothing" 拒绝方案里,有一段话我认为是整份文档最重要的部分(意译):

随着时间推移,这将成为上游实现最实质性的缺失功能。这会把高规模和云用户驱赶到 Kafka 替代品那里,它们的市场份额会增长。这会进一步分裂 Apache Kafka 对 Kafka 协议的控制,Kafka 可能完全失去对协议的主导权。这可能导致需要与分叉协调新功能、硬协议分叉激增,或者协议在某个标准组织下重新中心化。我们应该现在就采取措施来避免或延缓这个结果。

这不是技术文档的语气,这是战略文档的语气。

过去两三年发生了什么?WarpStream、AutoMQ、Confluent Freight、Redpanda Cloud Topics……一堆"Kafka 协议兼容"的产品,都在做同一件事:把活跃段扔进对象存储。它们证明了这个方向的市场价值,也证明了上游 Kafka 在这个方向上落后了。

更危险的是:当替代品足够多、份额足够大时,"Kafka 协议"就不再是 Apache Kafka 定义的了。 到那时候,Apache Kafka 想加个协议特性,得先问问几个大分叉答不答应。这就是 KIP 里说的"失去对协议的主导权"。

KIP-1150 的选择是:把创新收编回来。 用 Apache 2.0 许可证,让所有人都能免费用;让社区维护,减少对厂商的依赖;并且——这是关键——让后续的协议级优化(比如生产者机架感知)变得可能

最后这一条才是杀手锏。厂商的分叉可以实现 diskless,但它们改不了 Kafka 协议本身。而上游可以。PreferredProduceBrokers 这个 Metadata 字段就是明证——只有握着协议定义权的人,才能做这种改动。

我个人的判断:这一步走对了,但走晚了。 晚了大约两年。这两年里,市场教育已经被替代品完成了,用户心智里"云原生 Kafka = 对象存储"已经建立。上游现在做的事情,更多是收复失地而不是开疆拓土。

但收复失地本身就很有价值。对绝大多数企业用户来说,"上游原生支持"和"厂商分叉支持"之间的信任差距是巨大的。


十五、还有哪些没做完

KIP-1150 列出的后续工作(都还没有对应的 KIP,欢迎贡献):

  1. Topic 类型互转:classic ↔ diskless 双向转换。目前只能建新 topic 再迁数据,这是迁移路上最大的摩擦。
  2. Broker 角色专门化:把 Broker 按 produce / consume / coordination / compaction 分工,允许异构集群。这个方向一旦落地,Kafka 就真正变成了存算分离架构。
  3. 并行 Produce 处理:同时处理多个 Produce 请求,提升高延迟环境下的单生产者吞吐。跟 KIP-1269(可配置的 Broker 保留批次数)配套。
  4. Iceberg 格式:允许对静态 topic 数据做大规模并行处理。这条我最期待——它意味着 Kafka topic 可以直接被 Spark/Trino/Flink 当表读,"流表二象性"从概念变成物理现实。KIP 说这项工作会打开一个可插拔的存储接口,让日志格式层可以独立创新。
  5. 多区域 active-active:通过复制 topic 元数据实现自动故障转移。注意这里的措辞——只复制元数据。因为数据本来就在对象存储里,跨区域只需要同步"顺序"这个信息。这个设计的优雅程度让人眼前一亮。

另外还有一批相关 KIP 已经在关联中:

  • KIP-1279 集群镜像:diskless topic 可以通过引用对象存储中已有数据来高效镜像(不用重传数据,只传元数据引用)
  • KIP-1272 分层存储支持压实 topic:diskless topic 可以通过卸载到分层存储并在那里执行压实
  • KIP-1248 / KIP-1254 消费者直读分层存储:滞后的消费者可以直接从远程存储拉,减少 Broker 带宽压力

十六、总结:一句话记住每一层

我把整套设计压缩成一张表,方便你记:

层次一句话
动机跨 AZ 流量费占 Kafka TCO 的大头,必须干掉
核心手法数据平面无序(对象存储),控制平面定序(Coordinator)
存储格式WAL Segment 混装多分区,把 PUT 成本从 O(分区数) 降到 O(Broker 数)
写入路径任意 Broker 可写 → 缓冲 250ms/4MiB → 上传对象 → 提交坐标 → 分配 offset
延迟代价P50 ~500ms,P99 ~1-2s,是经典模式的 20-100 倍
协调器__diskless_metadata 多分区分片,SQLite 本地物化,推测性执行做流水线
副本语义in-sync 从"与 Leader 同步"变成"与 Coordinator 同步",Leader 也可能落后
GC 机制随机排列的所有者链式移交 + 定期孤儿文件扫描
成本优化引导生产者到同 AZ 的 DC gateway,把 N 次提交压到 1 次
选型原则延迟 SLA > 5s 且吞吐 > 50MB/s → 上;否则算清账再说

最后说点主观的。

我做过好几年消息中间件相关的活儿,Diskless 这套设计让我最欣赏的地方,不是它省了多少钱,而是它极其克制地控制了改动面

  • 消费者 API 一行不改
  • 生产者 API 一行不改(只需调参)
  • 事务、幂等、消费组、share group 全部保持语义
  • classic / tiered / diskless 三种 topic 在同一个集群里共存
  • 连 ISR 心跳都是复用老的 Fetch 协议掏空实现的

在一个已经跑了十几年、承载了无数生产系统的项目上做存储层重构,能做到"用户几乎无感",这才是真本事。 相比之下,那些动不动就"全新架构、不兼容升级"的方案,往往死在迁移路上。

当然,500ms 的延迟是个硬伤,而且是物理规律决定的硬伤,短期内没有解法。所以 Diskless 不会取代 classic topic,它们会长期共存——这也正是 KIP 设计成 per-topic 而非 per-cluster 的原因。

Kafka 从"一个消息队列",正式变成了"一个支持宽延迟谱系的流式引擎"。 这是 KIP-1150 原文里的说法,我觉得这个定位相当准确。

如果你的集群账单里跨 AZ 流量那一栏很扎眼,现在就可以开始盘:哪些 topic 的延迟 SLA 大于 5 秒? 那些就是你的第一批候选。


参考资料

  • KIP-1150: Diskless Topics(状态:Accepted)
  • KIP-1163: Diskless Core(状态:Under Discussion)
  • KIP-1164: Diskless Coordinator(状态:Under Discussion)
  • KAFKA-19161(JIRA 追踪)
  • KIP-405: Tiered Storage
  • KIP-630: Kafka Raft Snapshot
  • KIP-1269 / KIP-1272 / KIP-1279 / KIP-1248 / KIP-1254(关联 KIP)

本文中的成本计算基于公开单价的粗略估算,实际费用请以云厂商账单为准。KIP-1163 与 KIP-1164 仍在讨论中,最终实现的配置项名称和默认值可能与本文推演有出入,以正式发布的官方文档为准。

推荐文章

Elasticsearch 条件查询
2024-11-19 06:50:24 +0800 CST
Go中使用依赖注入的实用技巧
2024-11-19 00:24:20 +0800 CST
微信小程序开发资源汇总
2026-05-11 16:11:29 +0800 CST
Rust 中的所有权机制
2024-11-18 20:54:50 +0800 CST
Go语言中实现RSA加密与解密
2024-11-18 01:49:30 +0800 CST
程序员茄子在线接单