9082 字
45 分钟
Kafka 生产架构与治理:从 Topic、Partition 到容量、SLO 与跨集群灾备

Kafka 生产架构与治理:从 Topic、Partition 到容量、SLO 与跨集群灾备#

把 Kafka 底层机制转化为可计算、可观测、可演练、可持续演进的生产系统

阅读目标#

前六篇已经完成 Kafka 从抽象到底层、再到 Java 工程框架的完整拆解:

第一篇
Log → Partition → Offset → Consumer Position
第二篇
RecordBatch → LogSegment → Index → Page Cache → Data Path
第三篇
KafkaProducer.send()
→ RecordAccumulator
→ Sender
→ Partition Leader
第四篇
poll()
→ Consumer Group
→ Rebalance
→ Offset Commit
第五篇
Replication
→ ISR / HW / ELR
→ Idempotence
→ Transaction
第六篇
Spring Kafka
→ Listener Container
→ Retry / DLT
→ Transaction / Outbox

第七篇不再重点解释某个 API 或源码类,而是回答:

当你真正负责一个 Kafka 集群和数十个业务团队时,如何把前六篇的机制转化成架构决策、容量模型、SLO、变更流程和故障 Runbook?

本文用一个贯穿场景展开:

峰值写入:100,000 events/s
平均事件:1 KiB
同订单事件必须有序
数据保留:7 天
消费端到端延迟:P99 < 1 s
单 Broker 故障:不丢已确认事件
机房故障:RPO < 60 s,RTO < 30 min

完成本篇后,你应该能够:

  • 从业务语义和 SLO 推导 Kafka 拓扑,而不是先拍脑袋创建 Topic;
  • 制定 Topic 命名、归属、Retention、Schema 和生命周期规范;
  • 选择 Partition Key,并识别顺序、热点和数据倾斜风险;
  • 使用吞吐、消费能力和 Broker 分布三类约束估算 Partition 数;
  • 估算本地磁盘、复制、峰值和迁移期间的真实容量;
  • 设计生产 KRaft Controller Quorum 和 Broker 故障域;
  • 使用 ACL、Quota 和租户边界保护共享集群;
  • 建立 Produce、Consume、Replication、Storage、Coordinator 和 Controller 的统一 SLO;
  • 系统诊断 Consumer Lag,而不是看到 Lag 就增加 Consumer;
  • 安全执行 Partition Reassignment、Broker 扩容和下线;
  • 判断 Tiered Storage 适用边界和当前限制;
  • 设计 MirrorMaker 2 灾备拓扑,明确 RPO、RTO 和写入所有权;
  • 为 Topic 建立申请、发布、变更、回放和退役流程;
  • 把 Kafka 架构应用到真实的履约争议与 Multi-Agent 系统。

本文以 Apache Kafka 4.3.1 为版本基线。Kafka 4.3.x 只支持 KRaft;生产环境应将 Broker 与 Controller 角色分离,并通常使用 3 或 5 个 Controller。Kafka 4.3 还提供 Broker/Log Directory Cordon、Partition Size Monitoring 和 Tiered Storage 运维能力,但 Tiered Storage 仍要求用户提供 RemoteStorageManager 实现,并且当前不支持 Compacted Topic。


0. 架构题:10 万 TPS 的 Kafka 应该建多少 Partition#

假设收到需求:

订单事件峰值 100,000 events/s
每条平均 1 KiB
同订单严格有序
保留 7 天
P99 消费延迟 < 1 s
跨节点故障不丢已确认事件

很多方案会立刻回答:

创建 100 个 Partition
RF=3
acks=all

问题是:

为什么是 100?
单 Partition 在目标硬件和消息模型下能写多少?
单 Consumer 能处理多少?
100 个 Partition 是否足以分布到未来 Broker?
订单 Key 分布是否均匀?
Retention 需要多少磁盘?
消费者是否依赖全局顺序?
Schema 如何升级?
一个超级商家占 30% 流量怎么办?
故障和 Reassignment 时容量是否仍够?

Kafka 生产架构决策链

Kafka 架构应该遵循:

业务语义
SLO / Failure Budget
事件契约和 Key
Topic / Partition / Replica
容量与成本
观测、变更、恢复和治理

而不是:

先建 Topic
线上出问题
不断补参数

生产架构不是一个 server.properties,而是一组可以解释、测量、演练和回滚的决策。


1. 先定义 SLO,再讨论 Kafka 配置#

1.1 Kafka 平台至少有五类 SLO#

写入可用性#

Producer 请求成功率
Produce P99/P999 延迟
Delivery Timeout 比例

例如:

99.99% 的事件在 2 秒内得到最终发送结果

消费时效性#

Event Created Time
Business Side Effect Completed Time

定义端到端延迟,而不只是 Broker Fetch Latency:

E2E Processing P99 < 1 s

持久性#

明确:

在 RF=3、min ISR=2、acks=all 下,
允许哪些并发故障而不丢已确认事件?

不能笼统写“消息绝不丢失”。

恢复目标#

RPO:最多可接受多少未复制数据
RTO:集群或业务多久恢复

数据保留#

Retention = 7 天

它同时意味着:

所有 Consumer 必须在 7 天内完成消费或建立远程/外部归档能力。

Kafka 官方将 retention.ms 描述为消费者读取数据的 SLA 边界之一。

1.2 SLI 必须与用户体验关联#

只监控:

Broker CPU < 60%

不能代表用户正常。

业务真正关心:

订单创建后多快进入履约系统?
裁决事件多久触发执行?
积压是否会超过退货时效?

因此平台指标需要映射:

Kafka SLI业务意义
Produce Error事件可能未进入系统
Consumer Lag业务状态更新延迟
Rebalance Duration消费暂停窗口
Under Min ISR当前无法满足声明持久性
Offline Replica故障恢复能力下降
Disk UsageRetention 与可写性风险
Transaction LSO Lagread_committed 业务停滞

1.3 Failure Budget 决定冗余#

如果业务要求:

任意一台 Broker 故障仍可写

常见设计:

RF=3
min ISR=2
acks=all

但如果还要求:

整个可用区故障仍可写

需要继续确认:

三个 Replica 是否跨三个 AZ?
剩余两个 AZ 是否还能满足 min ISR?
网络延迟是否满足 Produce SLO?
Controller Quorum 是否跨故障域?

参数只有放进故障模型才有意义。


2. Topic 设计:围绕稳定业务事件,而不是消费者列表#

按业务域设计 Topic

2.1 Topic 是一个治理边界#

Topic 不只是一个字符串,它同时决定:

业务所有权
Partition 数
顺序范围
Schema 生命周期
Retention / Cleanup Policy
Replication Factor
ACL / Quota
监控和成本
回放与删除流程

因此 Topic 过粗和过细都会产生长期成本。

2.2 不要为每个 Consumer 创建 Topic#

错误:

order-for-risk
order-for-notification
order-for-analytics

新增消费者就要修改 Producer。

更合理:

Topic: order.lifecycle.v1
Group: fulfillment
Group: risk
Group: notification
Group: analytics

Kafka Consumer Group 已经提供独立消费视图。

2.3 Topic 粒度的三个常见模式#

一个业务域一个 Topic#

order.lifecycle.v1

内部包含多个 Event Type。

优点:

  • 同 Key 的生命周期事件容易保持同 Partition 顺序;
  • Consumer 订阅简单;
  • Topic 数可控。

缺点:

  • 不同 Event Type 共用 Retention、Partition 和 ACL;
  • Schema 管理更复杂;
  • 某类大事件可能影响整体。

一个事件类型一个 Topic#

order.created.v1
order.paid.v1
order.cancelled.v1

优点:

  • Schema、Retention、ACL 和容量独立;
  • 订阅精确。

缺点:

  • Topic 和 Partition 数快速膨胀;
  • 跨事件生命周期顺序难表达;
  • 运维和 Metadata 成本提高。

按数据用途拆 Topic#

例如:

order.lifecycle.v1
order.audit.v1
order.snapshot.compacted.v1

这是合理的,因为它们具有不同:

Cleanup Policy
Retention
Schema
访问权限
消费模式

2.4 推荐命名规范#

示例:

<domain>.<entity-or-stream>.<purpose>.v<major>

例如:

fulfillment.case-events.business.v1
fulfillment.case-state.snapshot.v1
fulfillment.audit.security.v1

命名中建议体现:

  • 业务域;
  • 数据语义;
  • 用途;
  • 不兼容 Major Version。

不要把环境放进 Topic 名:

prod-order-events

更推荐不同环境使用不同 Cluster 或 Namespace/ACL 边界。

2.5 deletecompact 是不同数据产品#

Delete Retention#

适合事件历史:

OrderCreated
OrderPaid
OrderShipped

到达时间或容量边界后删除旧 Segment。

Log Compaction#

适合每个 Key 的最新状态:

Key=order-1001
Value=current order state

Compaction 不等于“只保留一条”,也不保证立即删除旧值。

compact,delete#

同时保留最新 Key 状态,并限制更长期的历史范围。

2.6 Topic 设计评审清单#

创建前至少回答:

谁是业务 Owner?
Event 是事实还是命令?
Key 是什么?
需要什么顺序?
平均和最大消息大小?
Retention 多久?
Delete 还是 Compact?
谁能 Produce / Consume?
Schema 如何演进?
峰值吞吐?
消费者列表?
回放是否有业务风险?
何时可以删除?

3. Schema 演进:生产者和消费者之间的时间契约#

Kafka Broker 只保存字节。它不负责理解:

JSON 字段是否删除
Avro 类型是否兼容
Protobuf Field Number 是否复用
业务枚举是否改变语义

Schema 演进与兼容性

3.1 为什么滚动发布需要兼容#

发布过程中可能同时存在:

Producer v1
Producer v2
Consumer v1
Consumer v2
历史 Record v0

Schema 不能只保证“新代码能读新消息”。

3.2 常见兼容方向#

Backward Compatibility#

新 Consumer 能读取旧数据。

适合先升级 Consumer,再升级 Producer。

Forward Compatibility#

旧 Consumer 能读取新 Producer 的数据。

适合 Producer 可能先发布的场景。

Full Compatibility#

同时满足新读旧、旧读新。

代价更高,但适合大型共享事件流。

3.3 安全演进原则#

通常更安全:

新增可选字段
提供默认值
新增枚举时让旧消费者容错
保留字段编号

高风险:

删除必填字段
修改字段类型
改变字段语义但名称不变
重用 Protobuf Field Number
改变 Key Schema

3.4 Key Schema 比 Value 更敏感#

Key 变化可能直接影响:

Partition Mapping
局部顺序
Compaction Identity
Kafka Streams State Store
下游幂等键

因此 Key Schema 应比 Value 更稳定。

3.5 Schema Registry 不属于 Kafka Broker#

Apache Kafka 本身没有内置通用 Schema Registry。

企业通常选择:

  • Confluent Schema Registry;
  • Apicurio Registry;
  • 自建 Contract Registry;
  • Protobuf/Avro IDL + CI Compatibility Check。

真正重要的不是品牌,而是流程:

Schema Owner
Versioning
Compatibility Policy
CI Gate
Generated Model
Deprecation Window
Consumer Inventory

3.6 不要只管理语法兼容#

下面的变化可能 Schema 检查通过,但业务不兼容:

amount 单位从元改成分
status=SUCCESS 含义变化
timestamp 从事件时间改成处理时间
字段允许为空但旧逻辑未处理

因此还需要:

Semantic Contract Test
Golden Event
Consumer Contract Test
Shadow Consumption

4. Partition Key:顺序、并行和热点的统一决策#

Partition Key 与热点

4.1 Key 定义局部顺序域#

如果要求同订单有序:

OrderCreated
OrderPaid
OrderCancelled

应使用:

key = orderId

同一序列化 Key 被映射到同一个 Partition。

Kafka 只能保证:

同一 Partition Log 顺序

它不能保证:

同一个业务实体跨 Partition 全局有序

4.2 Key 也是负载分布函数#

假设:

merchant-A 占 75% 事件
key=merchantId

一个 Partition 可能承担绝大多数流量。

此时:

增加 Partition
增加 Consumer

都不能拆分已经进入同一 Key 顺序域的流量。

4.3 Key 候选比较#

Key顺序语义分布风险常见用途
orderId单订单顺序通常高基数、较均匀订单生命周期
userId单用户顺序超级用户热点用户行为
merchantId单商家顺序大商家强热点商家账务
caseId单案件顺序大案件热点较少履约争议
null无实体顺序Sticky Batch 较均匀遥测、无序日志

4.4 Hot Partition 的治理#

改变顺序域#

如果业务不需要商家全局顺序,可以改为:

merchantId + orderId

业务分桶#

merchantId + bucket

但必须接受:

同 Merchant 跨 Bucket 不再全局有序

隔离超级租户#

将大客户迁入独立 Topic 或 Cluster:

orders.standard.v1
orders.merchant-a.v1

上游限速#

Quota 或业务 Rate Limit 防止热点冲击整个 Broker。

拆分事件大小#

大 Value 会放大热点网络和磁盘压力。可以将大对象存对象存储,Kafka 只传引用和摘要。

4.5 必须监控 Key 分布#

Kafka Broker 不理解业务 Key 热度。

应用侧应统计:

Top-N Key Event Rate
Key Cardinality
Partition BytesIn
Partition RecordsIn
Partition Lag
Max/Median Partition Load

5. Partition 数量:建立可计算模型#

Partition 容量模型

Partition 数至少由三类约束决定。

5.1 写入吞吐约束#

设:

目标峰值写入 = T_write
单 Partition 在目标配置下压测吞吐 = C_partition_write

则:

P_write >= T_write / C_partition_write

注意单 Partition 容量不是网上固定数字,它受到:

消息大小
Batch Size
Compression
acks / RF / min ISR
Broker 磁盘与网络
TLS
Producer 数量
Follower 复制

必须在接近生产硬件和配置下压测。

5.2 消费并行度约束#

设:

目标消费速率 = T_consume
单 Consumer 实例可持续处理 = C_consumer

则:

P_consume >= T_consume / C_consumer

因为同 Group 有效 Consumer 数不超过 Partition 数。

如果消费逻辑包含数据库或外部 API,容量通常由下游而不是 Kafka Fetch 决定。

5.3 Broker 分布约束#

如果未来希望 Topic Leader 分布到 12 台 Broker,Partition 数不能只有 3。

P_distribution >= Desired Broker Spread

但 Partition 过多会增加:

Controller Metadata
文件与索引
Replica Fetcher 工作量
Leader Election 时间
Rebalance Assignment
恢复与 Reassignment 成本
客户端 Metadata

5.4 最终公式#

P_base = max(P_write, P_consume, P_distribution)
P_final = P_base × Growth Margin

Growth Margin 常根据增长速度和变更成本确定,而不是固定 2 倍。

5.5 贯穿案例#

假设压测结果:

单 Partition 可持续写 8 MiB/s
单 Consumer 可处理 2,000 events/s
目标 100,000 events/s
平均 1 KiB ≈ 97.7 MiB/s
希望至少分布到 12 Broker

写入约束:

97.7 / 8 ≈ 13 Partitions

消费约束:

100,000 / 2,000 = 50 Partitions

分布约束:

12 Partitions

所以:

P_base = 50

考虑 50% 增长余量:

P_final ≈ 75

工程上可以选择 72、80 或 96,但需要说明理由。

5.6 为什么不能随意后期加 Partition#

Kafka 支持增加 Partition,但不支持减少。

增加后:

  • hash(key) % partitionCount 映射可能变化;
  • 同 Key 后续事件可能进入新 Partition;
  • 历史数据不会自动重分布;
  • Consumer Group Rebalance;
  • Metadata 传播有延迟;
  • auto.offset.reset=latest 在特定窗口可能跳过新 Partition 的已有数据。

所以扩 Partition 是数据语义变更,不只是运维操作。

5.7 不修改内部 Topic Partition#

不要手动修改:

__consumer_offsets
__transaction_state
__share_group_state
__cluster_metadata

它们的 Partition 映射与 Coordinator/Metadata 逻辑相关。


6. 存储与 Retention 容量规划#

Kafka 存储容量模型

6.1 Primary Data 容量#

近似:

PrimaryBytes
≈ IngressBytesPerSec
× RetentionSeconds
× EffectiveCompressionRatio

如果 Producer 入口统计已经是压缩后的 Broker BytesIn,则不要重复乘压缩率。

6.2 集群物理容量#

ClusterBytes
≈ PrimaryBytes
× ReplicationFactor
× SafetyFactor

SafetyFactor 需要覆盖:

  • 峰值与增长;
  • Segment 和 Index;
  • Reassignment 临时双份;
  • Follower Catch-up;
  • Topic 增长不均;
  • 删除延迟;
  • 磁盘健康水位。

6.3 贯穿案例#

100,000 events/s
1 KiB/event
≈ 97.7 MiB/s
7 天 = 604,800 s

未压缩主数据约:

97.7 MiB/s × 604,800
≈ 56.4 TiB

假设有效压缩为 0.45:

Primary ≈ 25.4 TiB

RF=3:

Physical ≈ 76.2 TiB

再加入 1.35 Safety Factor:

≈ 102.9 TiB usable storage

这只是第一轮估算。还需要用实际 Broker BytesIn 和 Segment Size 验证。

6.4 retention.bytes 按 Partition 生效#

Kafka 官方说明:

Topic 总 Retention Bytes
≈ retention.bytes × Partition Count

如果设置:

retention.bytes=100 GiB
partitions=80

Topic 主副本逻辑容量上限约:

8,000 GiB

再乘 RF。

6.5 Retention 删除不是精确到秒#

数据只有在:

Segment Closed
Retention Check 扫描
满足删除条件
File Delete Delay

后才真正释放空间。

所以 7 天 Retention 不意味着所有 Record 在第 7 天整点立即消失。

6.6 磁盘水位#

不要把 Broker 设计到 95% 才告警。

需要预留:

Reassignment
Follower Recovery
流量突增
删除滞后
故障 Broker 数据迁移

常见做法是定义:

Warning 60%~70%
Expansion Trigger 70%~75%
Critical 80%~85%

具体阈值应结合扩容 Lead Time 和磁盘规模。

6.7 不同 Topic 不应统一 Retention#

数据类型建议策略
实时业务事件3~14 天,本地快速 Replay
审计事件更长 Retention 或外部归档
Compacted StateCompact,必要时附加 Delete
大数据回填独立 Topic 和临时 Retention
Retry Topic根据最大重试窗口设置
DLT根据人工处理 SLA 设置

7. 生产 KRaft 集群拓扑#

KRaft 生产拓扑

7.1 Controller 与 Broker 分离#

Kafka 4.3 KRaft 支持:

process.roles=broker
process.roles=controller
process.roles=broker,controller

Combined Mode 适合开发和小型环境,但官方不建议关键生产环境使用,因为:

  • Controller 无法独立扩缩容;
  • Broker 数据面压力影响控制面;
  • 滚动升级和故障隔离更差。

7.2 Controller 数量#

典型:

3 Controllers → 容忍 1 个 Controller 故障
5 Controllers → 容忍 2 个 Controller 故障

需要存活多数派:

2N + 1 Controllers
容忍 N 个并发故障

不是 Controller 越多越好。Quorum 越大,写 Metadata 需要更多协调。

7.3 Dynamic Controller Quorum#

Kafka 4.1 起支持动态 Controller Membership。

生产新集群可使用:

controller.quorum.bootstrap.servers=...

配合 kafka-metadata-quorum.sh 添加或移除 Controller。

Controller 变更必须:

  1. 新节点先格式化并启动;
  2. 等待 Metadata Log 追平;
  3. 执行 Add Controller;
  4. 删除时先 Remove,再 Shutdown。

7.4 Broker 与 Rack Awareness#

为 Broker 配置:

broker.rack=az-a

Replica Assignment 应跨 Rack/AZ。

RF=3 如果三个 Replica 都在同一 AZ:

单 Broker 故障安全
整个 AZ 故障不安全

7.5 Controller 故障域#

Controller Quorum 也应跨故障域。

如果三个 Controller 同一机架,机架故障会使整个 Kafka 无法进行 Metadata 变更和 Leader 管理。

7.6 节点资源分离#

Broker 重点资源#

Network
Local Disk Throughput / Capacity
Page Cache
CPU for Compression/TLS/Request

Controller 重点资源#

Stable Low-Latency Disk
Metadata Log
Memory for Cluster Metadata
Reliable Network

7.7 Java 版本#

Kafka 4.3 支持 Java 17、21 和 25。官方当前建议使用最新 LTS Java 25,以获得性能、效率和支持优势,但实际生产还需结合:

  • 组织 Java 基线;
  • 监控 Agent 兼容性;
  • GC 压测;
  • 安全补丁;
  • 客户端版本矩阵。

8. 多租户、安全与 Quota#

Kafka 多租户与 Quota

8.1 ACL 解决“能不能做”#

至少控制:

Topic READ / WRITE / CREATE / ALTER
Consumer Group READ
Transactional ID WRITE / DESCRIBE
Cluster ALTER / DESCRIBE

每个应用使用独立 Principal,不共享万能账号。

8.2 Quota 解决“能做多少”#

Kafka 支持按 User、Client ID 或组合设置:

producer_byte_rate
consumer_byte_rate
request_percentage
controller_mutation_rate

Byte Rate#

防止单租户占满网络。

Request Percentage#

限制该租户占用 Broker Request Handler CPU 时间。

在共享集群中,Request Quota 经常比纯带宽 Quota 更重要,因为大量小请求可能先耗尽 CPU。

Controller Mutation Rate#

限制 Topic Create/Delete/Alter 等控制面操作,避免租户通过高频 Admin 操作冲击 Controller。

8.3 Tenant Isolation 层级#

共享 Topic#

隔离最弱,不同租户事件进入同一 Partition 集合。

独立 Topic#

可独立 ACL、Quota、Retention 和 Partition。

独立 Cluster#

适用于:

  • 高合规数据;
  • 超大租户;
  • 独立升级窗口;
  • 强故障隔离;
  • 地域和数据主权要求。

8.4 安全不仅是 TLS#

Kafka 安全需要:

In-transit Encryption
Authentication
Authorization
Secret Rotation
Audit Log
PII Classification
Retention / Deletion
Consumer Export Control

尤其 DLT 和 Audit Topic 可能保存原始敏感 Value,不能因为“只是失败消息”就降低权限。


9. 生产可观测性:五个平面#

Kafka 可观测性五个平面

9.1 Client Plane#

Producer:

record-send-rate
record-error-rate
record-retry-rate
request-latency
record-queue-time
bufferpool-wait-time

Consumer:

records-consumed-rate
fetch-latency
records-lag-max
commit-latency
rebalance-rate
last-poll-seconds-ago

9.2 Broker Request Plane#

RequestsPerSec
TotalTimeMs
RequestQueueTimeMs
LocalTimeMs
RemoteTimeMs
ResponseQueueTimeMs
NetworkProcessorAvgIdlePercent
RequestHandlerAvgIdlePercent

判断瓶颈在:

排队
Broker 处理
Follower 等待
网络返回

9.3 Storage Plane#

Partition Size
RetentionSizeInPercent
Log Flush
Disk Latency
Disk Usage
Remote Copy Lag
Remote Fetch

Kafka 4.3 增加/强化 Partition Size Monitoring,允许直接观察单 Partition 大小和相对 Retention Limit 的比例。

9.4 Replication Plane#

UnderReplicatedPartitions
UnderMinIsrPartitionCount
OfflineReplicaCount
IsrShrinks/Expands
Replica Fetcher Lag
ReassigningPartitions
Reassignment Bytes In/Out

9.5 Control Plane#

KRaft Quorum Leader
Metadata Log End Offset
Follower Lag
Controller Event Queue
Consumer Group State
Transaction Coordinator

9.6 统一标签#

至少保留:

cluster
brokerId
rack
clientId
groupId
topic
partition
tenant
application

没有统一标签,无法把:

业务请求
→ Producer Client
→ Broker Partition
→ Consumer Group
→ 业务处理

串成一条链。

9.7 推荐 SLO Dashboard#

写入#

  • Produce P50/P95/P99;
  • Error / Retry;
  • Throttle;
  • Under Min ISR。

消费#

  • E2E Event Age;
  • Lag by Partition;
  • Processing P99;
  • Rebalance Duration。

集群#

  • Disk Forecast;
  • Network Saturation;
  • URP / Offline Replica;
  • Controller Quorum Health。

成本#

  • TiB retained by Topic;
  • Bytes In/Out by Tenant;
  • Partition Count by Owner;
  • Idle Topic / Consumer。

10. Consumer Lag:诊断,不是条件反射式扩容#

Consumer Lag 诊断树

10.1 基础模型#

Produce Rate = λ
Process Rate = μ

当:

λ > μ

Lag 持续增加。

但需要确定是:

所有 Partition
还是单个 Partition

10.2 输入突增#

检查:

  • 正常业务峰值;
  • 上游 Replay;
  • Producer Retry Storm;
  • 批处理任务;
  • 恶意或错误租户。

10.3 并行不足#

Partitions=12
Consumers=6

可增加 Consumer。

但:

Partitions=12
Consumers=12

再加 Consumer 无法提高 Partition 并行度。

10.4 Hot Partition#

如果:

P0 Lag=10M
P1..P11 Lag≈0

问题是 Key/流量倾斜,而不是 Group 总实例数。

10.5 处理变慢#

检查:

DB Connection Pool
SQL P99
HTTP Timeout
Thread Pool Queue
GC / CPU
Poison Record
Retry BackOff
DLT Publisher

10.6 Kafka Fetch 侧#

Broker Request Latency
Fetch Throttle
Leader Migration
Disk Read
Remote Tier Read
TLS CPU
Network Saturation

10.7 Lag 恢复时间#

如果当前积压:

Backlog = B
当前输入 = λ
扩容后处理 = μ

只有:

μ > λ

才会下降。

理论恢复时间:

T_recovery ≈ B / (μ - λ)

例如:

Backlog=10,000,000
输入=50,000/s
处理=70,000/s
T≈500 s

还需考虑 Rebalance 和下游限流。

10.8 Backpressure#

当下游无法扩容时,正确策略可能是:

  • pause/resume;
  • 上游限流;
  • 降低非关键事件;
  • 独立大租户;
  • 延迟 Retry Topic;
  • 扩展 Partition 和 Consumer;
  • 扩展数据库。

不能让 Kafka Lag 成为无上限缓冲区。Retention 到期后数据仍会消失。


11. 扩容、Reassignment 与 Broker 下线#

增加 Broker 不会自动把旧 Partition 均衡迁过去。

Partition Reassignment 流程

11.1 扩容触发条件#

不要等磁盘快满才扩容。

触发可以包括:

Disk Forecast 到达 Lead Time
Network P99 接近安全水位
Partition / Broker 过多
Page Cache Hit 下降
Reassignment Recovery 太慢
未来业务增长已确认

11.2 加入 Broker#

KRaft 新 Broker:

  1. 使用相同 Cluster ID 格式化;
  2. 配置唯一 node.id
  3. 配置 Rack;
  4. 启动并确认注册;
  5. 验证磁盘和 Network Listener。

11.3 生成 Reassignment Plan#

Kafka 工具支持:

--generate
--execute
--verify

但官方明确说明,该工具不会自动研究真实负载并生成完美均衡方案。

自动方案主要考虑副本分布,不理解:

Partition 热度
Tenant
磁盘速度差异
Leader Network
业务优先级

11.4 保存回滚 Plan#

执行前保存:

Current Assignment
Proposed Assignment

出现问题时可以将 Current Assignment 作为回滚输入。

11.5 设置复制 Throttle#

Terminal window
kafka-reassign-partitions.sh \
--bootstrap-server kafka:9092 \
--execute \
--reassignment-json-file plan.json \
--throttle 50000000

过高:

占满 Broker Network / Disk
影响 Producer 和 Consumer

过低:

迁移 Lag 不收敛

如果:

目标 Partition 持续写入速率 > 迁移 Throttle

Follower 永远追不上。

11.6 Verify 并移除 Throttle#

完成后执行 --verify,Kafka 工具会移除相关 Throttle 配置。

遗留 Throttle 是常见事故原因,会导致后续副本复制长期受限。

11.7 Broker Cordon 与 Decommission#

Kafka 4.3 支持通过:

cordoned.log.dirs=*

将 Broker 的 Log Directory 标记为 Cordon,避免继续分配新副本。

下线流程:

Cordon
→ 找出所有 Replica
→ 生成迁移计划
→ Execute + Verify
→ 确认 Broker 无 Partition
→ Shutdown
→ kafka-cluster.sh unregister

工具目前不会自动生成完整 Broker Decommission 计划,管理员仍需规划目标副本。

11.8 Leader Balance#

副本均衡不等于 Leader 均衡。

需要观察:

Leader Count by Broker
BytesIn/Out by Broker
Request Rate
Preferred Replica Imbalance

一个 Broker 即使磁盘不多,也可能承担大量 Leader 流量。


12. Tiered Storage:适合什么,不适合什么#

Kafka Tiered Storage

12.1 两层存储#

Local Tier
→ Broker 本地磁盘
→ Active 和近期 Segment
→ Tail Read / Page Cache
Remote Tier
→ S3 / HDFS 等
→ 已完成历史 Segment
→ Backfill / Recovery / Replay

12.2 为什么需要 Tiered Storage#

传统 Kafka:

Broker Storage Capacity
决定 Retention 上限

Tiered Storage 解耦:

Broker Compute / Local Hot Storage
Long-Term Retention Capacity

适合:

  • 长时间事件保留;
  • 历史 Replay;
  • Broker 扩缩容希望减少本地迁移数据;
  • 热数据少、冷数据多。

12.3 当前关键配置#

Broker:

remote.log.storage.system.enable=true
remote.log.storage.manager.class.name=...
remote.log.storage.manager.class.path=...

Topic:

remote.storage.enable=true
local.retention.ms=...
retention.ms=...
local.retention.bytes=...
retention.bytes=...

12.4 Apache Kafka 不提供生产 RemoteStorageManager#

Kafka 提供 SPI,但当前不提供开箱即用的生产 RemoteStorageManager

必须选择或实现插件,并验证:

对象存储一致性
Multipart Upload
Retry / Idempotence
Credential Rotation
Encryption
Metrics
Delete Semantics
Upgrade Compatibility

12.5 当前限制#

Kafka 4.3 Tiered Storage 文档明确列出:

不支持 Compacted Topic
关闭 Broker 级 Tiered Storage 前必须先处理所有启用 Topic
某些旧 Segment 缺失 Producer Snapshot 不支持

因此不能把所有 Topic 都直接开启远程层。

12.6 性能边界#

实时 Consumer 通常读 Tail,本地 Page Cache 仍是主路径。

历史 Replay 访问远程层时:

  • 首字节延迟更高;
  • 受对象存储吞吐和请求限制;
  • Remote Cache 可能被击穿;
  • 大规模回放会影响生产流量。

需要为 Backfill 建立:

独立 Consumer Group
速率限制
时间窗口
优先级
成本预算

12.7 Tiered Storage 不等于灾备#

远程 Segment 不能替代:

跨集群复制
Metadata 灾备
Consumer Offset 灾备
应用 Failover

它主要解决存储层级,不是完整 DR。


13. 跨集群、MirrorMaker 2 与灾备#

多集群灾备模式

13.1 不建议一个 Kafka Cluster 跨超远地域#

Kafka Replica 协议对延迟敏感。跨大洲或高延迟网络组成一个 Quorum/Replica 集群,会放大:

Produce ACK Latency
Replica Lag
Controller Quorum 风险
网络分区影响

常见模式是:

每个地域本地 Kafka Cluster
应用只访问本地 Cluster
通过 MirrorMaker 2 跨集群复制

13.2 MirrorMaker 2 组件#

基于 Kafka Connect,主要包括:

MirrorSourceConnector
MirrorCheckpointConnector
MirrorHeartbeatConnector

用于复制:

  • Topic Data;
  • Consumer Group Checkpoint;
  • Heartbeat/Connectivity;
  • 部分配置和 ACL(取决于配置)。

13.3 Active-Passive#

Cluster A Active
Cluster B Warm Standby
A → B Mirror

优点:

  • 写入所有权清晰;
  • 冲突少;
  • Failover 流程可控。

需要设计:

Producer Endpoint 切换
Consumer Offset Translation
Topic Config 同步
Schema Registry 同步
Transactional ID 冲突
回切策略

13.4 RPO#

取决于:

Mirror Consumer Lag
Source/Target 可用性
Network Bandwidth
Checkpoint Frequency

RPO 不能只写“有 MirrorMaker 所以零丢失”。

13.5 RTO#

包含:

确认主集群不可恢复
启用 DR Producer
切换 DNS / Config
恢复 Consumer Group
校验重复和缺口
业务人工批准

13.6 Active-Active#

两地同时写入更复杂:

同一 Key 在两地修改
事件冲突
循环复制
Topic 名重写
重复事件
无法形成跨集群全局顺序

需要业务层定义:

  • 单 Key Home Region;
  • Conflict Resolution;
  • Global Event ID;
  • Loop Prevention;
  • 最终一致性规则。

13.7 DR 必须演练#

至少季度演练:

停止 Primary Produce
确认 Mirror Lag
切换 Producer
恢复 Consumer
验证 Offset
校验数据缺口和重复
执行回切
记录实际 RPO/RTO

纸面架构不能证明灾备有效。


14. Topic 全生命周期治理#

Topic 生命周期治理

14.1 Topic Catalog#

每个 Topic 至少登记:

Topic Name
Business Owner
Technical Owner
Data Classification
Event Description
Key Schema
Value Schema
Partition Count
RF / min ISR
Retention / Cleanup Policy
ACL / Quota
Producer List
Consumer Group List
SLO
Cost Center
Replay Procedure
Deletion Procedure

14.2 创建审批#

自动检查:

  • 命名规范;
  • Partition 上限;
  • RF;
  • min ISR;
  • Retention;
  • Schema Compatibility;
  • Owner;
  • ACL;
  • 容量预算。

不要允许业务随意开启 Auto Topic Creation。

14.3 变更审批#

高风险变更:

增加 Partition
降低 RF / min ISR
修改 Retention
Delete → Compact
开启 Tiered Storage
修改 ACL
大规模 Offset Reset
Reassignment

每次需要:

Impact Analysis
Rollback Plan
Observation Window
Change Owner
Audit Record

14.4 Replay 治理#

Replay 会重新触发:

数据库更新
HTTP 调用
通知
Agent Tool
计费

必须有:

  • Replay Group;
  • 时间/Offset 范围;
  • 速率限制;
  • 业务幂等确认;
  • DLT 策略;
  • 审批和审计。

14.5 Topic 退役#

流程:

停止 Producer
确认无新写入
盘点 Consumer
下线 Consumer
导出/归档必要数据
等待安全窗口
删除 ACL / Schema
删除 Topic
更新 Catalog

14.6 防止 Topic 墓地#

定期扫描:

无 Producer
无活跃 Consumer
BytesIn=0
长期无读取
Owner 已离职
Retention 无限
Partition 过多

无人负责的 Topic 是平台长期成本和安全风险。


15. 案例:履约争议 Multi-Agent 平台#

下面把整篇设计应用到你的“履约争议游园会”。

履约争议平台 Kafka 架构

15.1 核心事件流#

Topic: fulfillment.case-events.business.v1
Key: caseId

事件:

CaseCreated
EvidenceSubmitted
DossierBuilt
DisputeRouted
StatementSubmitted
JudgmentDrafted
JudgmentApproved
ExecutionStarted
ExecutionCompleted

使用 caseId 保证单 Case 生命周期局部有序。

15.2 Consumer Groups#

evidence-agent
policy-agent
adjudication-runtime
execution-service
audit-analytics
notification

每个 Group 独立消费同一事实流。

15.3 为什么不为四个 Agent 创建四份 Producer Topic#

Agent 是消费者角色,可能随产品迭代变化。

稳定的是:

Case 事实

所以 Producer 不应该知道有哪些 Agent。

15.4 状态快照 Topic#

fulfillment.case-state.snapshot.v1
cleanup.policy=compact
key=caseId

用于保存最新 Case State,减少新服务从完整历史重建的时间。

注意当前 Tiered Storage 不支持 Compacted Topic,因此该 Topic 不能直接采用远程分层存储。

15.5 审计 Topic#

fulfillment.audit.security.v1

特点:

  • 更长 Retention;
  • 严格 ACL;
  • 不允许普通业务 Consumer;
  • 记录 Agent Decision、Human Review 和 Tool Call;
  • 可能归档到不可变存储。

15.6 Agent Tool 副作用#

Execution Agent 可能调用:

退款 API
库存 API
优惠券 API
通知 API

必须使用:

caseId + actionId 作为 Idempotency Key
Action State Machine
Outbox / Inbox
结果查询
人工补偿

不能因为事件来自 Kafka 就认为 Tool Call Exactly Once。

15.7 Hot Key#

大型案件可能包含大量证据和评议事件。

治理:

  • 大 Value 放对象存储;
  • Kafka 传 Object URI、Hash、Metadata;
  • 对单 Case 设置事件速率限制;
  • Evidence Chunk 可使用 caseId + evidenceId 子流,但要明确顺序变化;
  • Judge State 仍按 caseId 串行更新。

15.8 PII 和数据删除#

争议证据可能包含:

姓名
地址
电话
订单信息
聊天记录
图片

设计:

  • Value 加密或 Tokenization;
  • Topic ACL 最小权限;
  • 大文件不直接进 Kafka;
  • Retention 与法律要求一致;
  • 删除请求需要覆盖 Kafka、对象存储、索引和下游数据仓库。

15.9 SLO#

示例:

CaseCreated → Evidence Agent 收到 P99 < 500 ms
JudgmentApproved → ExecutionStarted P99 < 2 s
核心事件 Produce Success > 99.99%
DLT 未处理时间 < 30 min
单 AZ 故障 RTO < 5 min

16. 可运行实验:Topic、Quota、Lag 与 Reassignment#

以下实验用于学习和演练,不是完整生产部署模板。

16.1 三 Broker KRaft Compose#

x-kafka-common: &kafka-common
image: apache/kafka:4.3.1
environment: &kafka-env
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
KAFKA_DEFAULT_REPLICATION_FACTOR: 3
KAFKA_MIN_INSYNC_REPLICAS: 2
services:
kafka-1:
<<: *kafka-common
hostname: kafka-1
ports: ["19092:19092"]
environment:
<<: *kafka-env
KAFKA_NODE_ID: 1
KAFKA_LISTENERS: INTERNAL://:9092,EXTERNAL://:19092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-1:9092,EXTERNAL://localhost:19092
kafka-2:
<<: *kafka-common
hostname: kafka-2
ports: ["29092:29092"]
environment:
<<: *kafka-env
KAFKA_NODE_ID: 2
KAFKA_LISTENERS: INTERNAL://:9092,EXTERNAL://:29092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-2:9092,EXTERNAL://localhost:29092
kafka-3:
<<: *kafka-common
hostname: kafka-3
ports: ["39092:39092"]
environment:
<<: *kafka-env
KAFKA_NODE_ID: 3
KAFKA_LISTENERS: INTERNAL://:9092,EXTERNAL://:39092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-3:9092,EXTERNAL://localhost:39092

Combined Mode 只用于本地实验。关键生产环境应分离 Controller 和 Broker。

16.2 创建生产风格 Topic#

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka-1:9092 \
--create \
--topic fulfillment.case-events.business.v1 \
--partitions 12 \
--replication-factor 3 \
--config min.insync.replicas=2 \
--config cleanup.policy=delete \
--config retention.ms=604800000 \
--config segment.bytes=1073741824

查看:

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka-1:9092 \
--describe \
--topic fulfillment.case-events.business.v1

16.3 创建 Compacted State Topic#

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka-1:9092 \
--create \
--topic fulfillment.case-state.snapshot.v1 \
--partitions 12 \
--replication-factor 3 \
--config min.insync.replicas=2 \
--config cleanup.policy=compact \
--config min.cleanable.dirty.ratio=0.5

16.4 查看 Topic 配置#

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-configs.sh \
--bootstrap-server kafka-1:9092 \
--entity-type topics \
--entity-name fulfillment.case-events.business.v1 \
--describe

16.5 制造 Consumer Lag#

生产大量数据:

Terminal window
python3 - <<'PY' | docker exec -i kafka-1 \
/opt/kafka/bin/kafka-console-producer.sh \
--bootstrap-server kafka-1:9092 \
--topic fulfillment.case-events.business.v1 \
--property parse.key=true \
--property key.separator=:
for i in range(200000):
print(f"case-{i % 5000}:event-{i}")
PY

启动慢 Consumer:

Terminal window
docker exec -it kafka-1 \
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka-1:9092 \
--topic fulfillment.case-events.business.v1 \
--group slow-lab

查看 Lag:

Terminal window
docker exec kafka-1 \
/opt/kafka/bin/kafka-consumer-groups.sh \
--bootstrap-server kafka-1:9092 \
--describe \
--group slow-lab

重点比较每个 Partition 的 Lag,而不是只看总数。

16.6 设置 Client Quota#

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-configs.sh \
--bootstrap-server kafka-1:9092 \
--alter \
--entity-type clients \
--entity-name noisy-producer \
--add-config producer_byte_rate=1048576,request_percentage=50

查看:

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-configs.sh \
--bootstrap-server kafka-1:9092 \
--describe \
--entity-type clients \
--entity-name noisy-producer

使用相同 client.id 发送,观察 Produce Throttle 指标和延迟。

16.7 Reassignment Plan#

创建 Topic 清单:

{
"topics": [
{"topic": "fulfillment.case-events.business.v1"}
],
"version": 1
}

生成:

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-reassign-partitions.sh \
--bootstrap-server kafka-1:9092 \
--topics-to-move-json-file /tmp/topics.json \
--broker-list "1,2,3" \
--generate

执行自定义 Plan:

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-reassign-partitions.sh \
--bootstrap-server kafka-1:9092 \
--reassignment-json-file /tmp/plan.json \
--execute \
--throttle 50000000

验证:

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-reassign-partitions.sh \
--bootstrap-server kafka-1:9092 \
--reassignment-json-file /tmp/plan.json \
--verify

16.8 Cordon Broker Log Directory#

Terminal window
docker exec kafka-1 /opt/kafka/bin/kafka-configs.sh \
--bootstrap-server kafka-1:9092 \
--alter \
--add-config cordoned.log.dirs="*" \
--entity-type brokers \
--entity-name 3

注意:Cordon 不会自动迁走已有 Replica,还必须执行 Reassignment。

16.9 实验验收#

完成后应能解释:

  • Topic 配置是 Broker Default 还是 Topic Override;
  • 12 个 Partition 如何分布到三 Broker;
  • Consumer Lag 是否集中在某个 Partition;
  • Quota 如何体现为 Client Throttle;
  • Reassignment 为什么需要保存旧 Plan;
  • Throttle 过低为什么不收敛;
  • Cordon 与真正 Decommission 的区别。

17. Admin API 与源码阅读地图#

第七篇的源码重点不再是单条 Record,而是控制面和运维状态。

17.1 AdminClient#

阅读:

Admin
KafkaAdminClient
AdminMetadataManager
AdminApiDriver

关注:

  • Topic Create/Alter/Delete;
  • Config Describe/Alter;
  • Consumer Group Describe;
  • ListOffsets;
  • Replica Reassignment;
  • Feature Version;
  • Broker Unregister。

17.2 KRaft Controller#

QuorumController
ReplicationControlManager
ClusterControlManager
ConfigurationControlManager
FeatureControlManager

关注:

  • Topic/Partition Metadata;
  • Replica Assignment;
  • Leader Election;
  • ISR/ELR 更新;
  • Broker Registration;
  • Cordon Log Dirs;
  • Metadata Log Record。

17.3 Log Manager 与 Retention#

LogManager
UnifiedLog
LocalLog
LogCleaner
RemoteLogManager

关注:

  • Retention Check;
  • Segment Delete;
  • Compaction;
  • Partition Size Metrics;
  • Remote Segment Copy/Delete;
  • Log Directory Failure。

17.4 Quota#

ClientQuotaManager
ClientQuotaCallback
ControllerMutationQuota
ReplicationQuotaManager

关注:

  • User/Client 匹配优先级;
  • Throttle 时间计算;
  • Request Percentage;
  • Reassignment Replication Quota。

17.5 MirrorMaker 2#

MirrorSourceConnector
MirrorCheckpointConnector
MirrorHeartbeatConnector
RemoteClusterUtils

重点不是逐行读,而是画出:

Source Topic
→ Source Consumer
→ Target Producer
→ Remote Topic
Source Group Offset
→ Checkpoint Topic
→ Offset Translation

17.6 推荐阅读顺序#

官方 Basic Operations
→ Admin API
→ KRaft Metadata
→ Reassignment
→ Retention / LogCleaner
→ Quota
→ Tiered Storage
→ MirrorMaker 2

18. 常见错误认知#

18.1 “Partition 越多吞吐越高”#

过多 Partition 会提高 Metadata、文件、Rebalance、选举和迁移成本。

18.2 “加 Broker 后会自动均衡数据”#

Kafka 不会自动移动旧 Replica,需要 Reassignment。

18.3 “Consumer Lag 高就加 Consumer”#

如果 Partition 已全部占用或是单 Hot Partition,加 Consumer 无效。

18.4 “Retention 7 天就是第 7 天准时删除”#

删除以 Segment 和周期检查为单位。

18.5 “Schema Registry 能保证业务兼容”#

它主要验证结构兼容,不能自动验证字段语义。

18.6 “Tiered Storage 等于备份”#

它不提供完整跨集群 DR 和应用 Failover。

18.7 “MirrorMaker 2 可以自动 Active-Active”#

工具复制数据,不解决业务冲突和全局顺序。

18.8 “RF=3,所以可以容忍三台 Broker 故障”#

RF=3 最多在数据意义上丢失两个副本,但是否仍可读写取决于 Leader、ISR、min ISR、Rack 和同时故障位置。

18.9 “共享 Kafka 只要 ACL 就够”#

没有 Quota,一个合法租户仍能耗尽 Broker CPU、网络和控制面。


19. 场景题#

19.1 Topic 磁盘增长比估算快一倍#

检查:

估算是否使用未压缩还是压缩流量?
实际消息是否变大?
Partition Count 是否增加?
RF 是否变化?
Retention 是否被覆盖?
Segment 是否无法及时删除?
Consumer Replay 是否产生额外 Topic?
Reassignment 是否临时双份?

19.2 新增 Broker 后 CPU 很低,但旧 Broker 仍满#

原因:旧 Partition 未迁移。

执行 Reassignment,并同时检查 Leader Distribution。

19.3 一个 Group Lag 高,其他 Group 正常#

说明 Kafka Topic 本身可能正常,优先检查该 Group:

Processing
Commit
Rebalance
Consumer Count
Downstream
Hot Partition

19.4 所有 Group 同时 Fetch 慢#

优先检查:

Broker Network
Disk
Request Queue
Quota
Leader Concentration
TLS CPU
Controller / Leader Election

19.5 扩 Partition 后同订单事件乱序#

默认 Hash 映射可能变化,同 Key 新事件进入新 Partition。

需要:

  • 发布窗口;
  • Producer Metadata 收敛;
  • 旧流排空;
  • 自定义稳定分区策略;
  • 新 Topic 迁移;
  • 接受短期双流并做版本路由。

19.6 DR 切换后重复处理#

原因可能:

  • Checkpoint Lag;
  • Target Group Offset 不精确;
  • Primary 最后事件已处理但未镜像 Offset;
  • Producer 双写窗口;
  • 应用无幂等。

DR 必须默认可能重复,并通过业务 Idempotency Key 保护。


20. 最终验收清单#

20.1 Topic#

可以说明:

为什么按这个业务域拆?
为什么不是每 Event Type 一个 Topic?
Cleanup Policy 是什么?
Owner 和 Consumer 是谁?

20.2 Partition#

能够计算:

写入约束
消费约束
Broker 分布约束
Growth Margin

并说明 Partition 扩容对 Key 顺序的影响。

20.3 Key#

可以解释:

顺序域
Key Cardinality
Hot Key
Compaction Identity
业务幂等

20.4 Storage#

能够估算:

Ingress × Retention × Compression × RF × Safety Factor

并说明 Reassignment、Segment 删除和磁盘水位余量。

20.5 KRaft#

能够设计:

3/5 Controller Quorum
Broker/Controller 分离
Rack Awareness
Dynamic Controller Membership
Broker Registration / Unregister

20.6 Multi-Tenancy#

能够区分:

Authentication
ACL
Byte Quota
Request Quota
Controller Mutation Quota
独立 Topic / Cluster 隔离

20.7 Observability#

至少建立:

Client
Broker Request
Storage
Replication
Coordinator / Controller

五个平面 Dashboard。

20.8 Lag#

面对 Lag 先回答:

所有 Partition 还是单热点?
λ 和 μ 分别是多少?
Group 是否 Stable?
Consumer 是否有空闲 Partition?
Downstream 是否饱和?
预计恢复时间?

20.9 Reassignment#

能够执行:

Generate
Review
Save Rollback
Execute with Throttle
Monitor Replica Lag
Verify
Remove Throttle

20.10 DR#

明确:

RPO
RTO
Primary Writer
Failover Trigger
Offset Translation
Duplicate Handling
Failback
Drill Frequency

20.11 最终口述验收#

不看文章,用 60 分钟完成架构设计:

一个事件平台峰值 10 万条每秒,每条约 1 KiB,要求同订单有序、保留 7 天、P99 端到端延迟低于 1 秒,单 Broker 故障不丢已确认数据,并在另一地域建设 60 秒 RPO 的灾备。请从 Event、Topic、Schema、Key、Partition Count、RF/minISR、磁盘、KRaft、Quota、Lag SLO、Reassignment、Tiered Storage 和 MirrorMaker 2 给出完整方案,并解释每个选择的代价。

如果回答仍然只有:

创建 Kafka 集群,设置 100 个 Partition,RF=3,开启监控。

说明还没有达到 Kafka 生产架构能力。


21. 全文收束:Kafka 工程专家的能力闭环#

七篇文章最终形成一条完整链:

Distributed Log Mental Model
Storage Engine / Data Path
Producer Runtime
Consumer Coordination
Replication / Transaction
Spring Kafka Engineering
Production Architecture / Governance

Kafka 工程专家不是掌握最多参数的人。

他应该能够在下面几个层次之间往返:

业务语义#

什么是事件?
谁拥有它?
顺序边界是什么?
什么才算处理成功?

分布式机制#

Partition
Offset
Replica
ISR
Coordinator
Transaction

工程实现#

Producer / Consumer
Spring Listener
Retry / DLT
Inbox / Outbox

生产治理#

SLO
Capacity
Security
Quota
Metrics
Runbook
DR
Cost
Lifecycle

最终的核心原则是:

每一个 Kafka 参数,都应该能追溯到一个业务语义、容量约束或故障模型;每一个架构承诺,都必须有指标、演练和恢复流程证明它成立。

到这里,Kafka 七篇核心专题完成闭环。

后续可以作为独立扩展专题继续研究:

Kafka Streams 与 Stateful Processing
Kafka Connect 与 CDC
Share Groups
Schema Registry 深入
Tiered Storage 实现
MirrorMaker 2 与跨地域一致性
KRaft Controller 源码
Kafka 性能压测与 Benchmark 方法论

参考资料#

  1. Apache Kafka 4.3.x Documentation
  2. Apache Kafka 4.3 — Basic Kafka Operations
  3. Apache Kafka 4.3 — Monitoring
  4. Apache Kafka 4.3 — KRaft
  5. Apache Kafka 4.3 — Topic Configs
  6. Apache Kafka 4.3 — Broker Configs
  7. Apache Kafka 4.3 — Tiered Storage
  8. Apache Kafka 4.3 — MirrorMaker Configs
  9. Apache Kafka 4.3 — Security Overview
  10. Apache Kafka 4.3 — Java Version
  11. Apache Kafka 4.3 — Eligible Leader Replicas
  12. Apache Kafka Source Repository
Kafka 生产架构与治理:从 Topic、Partition 到容量、SLO 与跨集群灾备
https://jupiter-ws.cn/posts/backend/kafka/07_kafka_production_architecture_governance/
作者
Jupiter
发布于
2026-07-18
许可协议
CC BY-NC-SA 4.0