Kafka 概览
必要且要理解消息语义的边界。topic / consumer 的样板代码可以让 AI 写,但分区、副本、消费组、幂等与顺序性、Exactly-once 到底到哪一层,这些决定了消息会不会丢或重复——AI 不会告诉你它生成的配置其实有丢消息风险。代码不必死记,用时让 AI 生成并人工核对;但你要把消息语义边界学到「能审 AI 的配置」的程度,否则线上丢数据或重复消费你查不出根因。
如果你写过后端服务,一定遇到过这样的场景:订单服务下单成功后,要通知库存扣减、发短信、更新报表、同步数仓——直接用 RPC 一个个调过去,调用链越来越长,任何一个下游挂掉都会拖垮下单接口。消息系统就是为了解开这种耦合而生的,而 Kafka 是其中吞吐量最高、生态最成熟的一个。
1. 为什么需要消息系统
同步调用架构有三个根本问题:
- 耦合:订单服务必须知道所有下游的存在。新增一个"风控分析"消费方,就得改订单服务的代码。
- 削峰:大促时下单 QPS 冲到 10 万,但库存服务只能处理 2 万,同步调用直接把库存打挂。
- 可靠性:下游短暂不可用时,请求直接失败,没有缓冲和重试的余地。
引入消息系统后,订单服务只做一件事——把"订单已创建"这个事件写进消息系统,然后立刻返回。下游各自按自己的节奏拉取消费:
同步调用(紧耦合):
订单服务 ──> 库存服务
├────> 短信服务
├────> 报表服务
└────> 数仓同步 任何一个慢/挂 => 下单变慢/失败
引入 Kafka(解耦):
订单服务 ──> [ Kafka: orders topic ] <── 库存服务 (按自己节奏拉)
<── 短信服务
<── 报表服务
<── 数仓同步 新增消费方无需改上游2. Kafka 的定位:分布式提交日志
很多人把 Kafka 当"消息队列",但它的本质是分布式、可分区、可复制的提交日志(commit log)。理解这一点是理解 Kafka 一切设计的钥匙。
- 日志(log):消息只会追加到文件末尾,写入后不可修改。顺序写磁盘的速度接近内存随机写,这是 Kafka 高吞吐的根基。
- 消费即读取偏移量:消息被消费后不会删除(由保留策略统一清理)。每个消费者只记录自己读到了第几条(offset),所以同一份数据可以被无数个消费组独立、重复地读取。
- 分区(partition):一个 topic 拆成多个分区,分布在不同机器上,写入和消费都能水平扩展。
- 复制(replication):每个分区有多个副本,容忍机器宕机。
Topic: orders (3 个分区)
partition-0: [msg0][msg1][msg2][msg3] ──> 追加写
partition-1: [msg0][msg1][msg2] ──> 追加写
partition-2: [msg0][msg1][msg2][msg3][msg4]
消费组 A 读到 p0 的 offset=2
消费组 B 读到 p0 的 offset=4 两者互不影响传统队列(如 RabbitMQ 经典模式)的消息被消费后就删除,"回放历史"很困难;而 Kafka 的日志模型天然支持回放、多订阅方、流处理,这也是它成为"事件流平台"而非单纯队列的原因。
3. 典型应用场景
| 场景 | 说明 | 例子 |
|---|---|---|
| 异步解耦 | 上游发事件,下游各自消费 | 订单创建后通知多个子系统 |
| 削峰填谷 | 高峰流量先落 Kafka,下游匀速消费 | 秒杀、大促 |
| 日志采集 | 应用/访问日志统一汇聚 | Filebeat → Kafka → ES |
| 流处理 | 实时计算、窗口聚合 | Kafka Streams / Flink 实时大屏 |
| 变更捕获 CDC | 数据库 binlog 流入 Kafka | Debezium 同步 MySQL 到数仓 |
| 事件溯源 | 以事件日志作为事实来源 | 账务系统、审计 |
4. 与 RabbitMQ / RocketMQ / Pulsar 对比
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 模型 | 分布式日志 | 传统队列(AMQP) | 日志 + 队列混合 | 日志,存算分离 |
| 单机吞吐 | 百万级 msg/s | 万级 | 十万级 | 十万级 |
| 消息回放 | 原生支持 | 困难 | 支持 | 支持 |
| 延迟消息 | 不原生支持 | 插件支持 | 原生支持 | 原生支持 |
| 路由灵活性 | 弱(按 topic/分区) | 强(exchange 路由) | 中(tag 过滤) | 中 |
| 生态 | 最成熟(Connect/Streams/Flink) | 一般 | 国内电商场景多 | 较新 |
| 运维复杂度 | 中(KRaft 后降低) | 低 | 中 | 高(BookKeeper) |
选型建议:
- 大数据量、日志流、流处理、需要回放 → Kafka
- 业务消息量不大、需要复杂路由/延迟队列/优先级 → RabbitMQ
- 电商交易场景、需要事务消息与定时消息 → RocketMQ
- 多租户、跨地域复制、存算分离诉求强烈 → Pulsar
Kafka 不适合做任务队列(没有单条消息的 ACK/重投递语义)、不原生支持延迟消息和优先级队列、单条消息默认最大 1MB(可调但不建议放大很多)。如果你的需求是"给某个 worker 派发任务并确认完成",RabbitMQ 或专门的任务队列更合适。
5. 版本与运行模式:KRaft 取代 ZooKeeper
Kafka 早期依赖 ZooKeeper 存储元数据(broker 注册、topic 配置、controller 选举),带来额外的运维负担和分区数扩展瓶颈。从 2.8 开始引入 KRaft(Kafka Raft)模式,用内置的 Raft 共识协议管理元数据;3.3 起 KRaft 生产可用;从 4.0 起彻底移除 ZooKeeper。本课程以 Kafka 3.7+ 与 KRaft 模式为基准。
ZooKeeper 模式(历史): KRaft 模式(现在):
+-----------+ +--------------------+
| ZooKeeper | <- 元数据/选举 | Controller (Raft) | 内置
+-----+-----+ +--------------------+
+-----+-----+ | Broker |
| Brokers | +--------------------+
+-----------+ 单一系统, 部署更简单,
两套系统分别运维 支持百万级分区先跑起来(第 3、4 章的安装与 CLI),再学 Producer/Consumer 客户端(5-10 章),然后深入原理(11、12 章副本与存储),最后是生态与运维(13-19 章),第 20 章用一个订单流水线项目串起全部知识。原理部分不要跳过——Kafka 的大多数线上事故都源于对副本同步与 Rebalance 机制的误解。
6. 一条消息的旅程(预览)
用一条订单消息串起后面章节会展开的所有概念:
Producer Kafka 集群 Consumer Group
-------- --------- --------------
序列化 key/value topic: orders poll() 拉取
| partition 由 key 哈希决定 |
v | v
按分区攒批 (batch) ---> leader 副本追加写日志 反序列化
| | |
v v v
Sender 线程发送 follower 副本拉取同步 处理业务, 提交 offset
|
v
ISR 全部确认 => acks=all 返回成功- 消息发到哪个分区?→ 第 10 章分区策略
- 怎么保证不丢?→ 第 6 章 acks 与第 11 章 ISR
- 消费者挂了怎么办?→ 第 9 章 Rebalance
- offset 存在哪?→ 第 8 章
小结
- 消息系统解决解耦、削峰、可靠缓冲三大问题
- Kafka 的本质是分布式提交日志:追加写、按 offset 读、消费不删除、支持回放
- 高吞吐来自顺序写磁盘、批量发送、零拷贝(后续章节展开)
- 3.7+ 时代用 KRaft 模式,不再需要 ZooKeeper
- 选型:大流量与流处理选 Kafka,复杂路由选 RabbitMQ,事务/延迟消息选 RocketMQ
- 列出你当前系统中三个可以用消息系统解耦的调用链,并说明用 Kafka 还是 RabbitMQ 更合适、为什么。
- 思考:Kafka "消费后不删除消息"的设计,对实现"新上线一个数据分析服务、需要重算过去 7 天数据"这类需求有什么帮助?