Learn
Kafka/01-introduction

Kafka 概览

💡🤖 AI 时代,还要学 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 流入 KafkaDebezium 同步 MySQL 到数仓
事件溯源以事件日志作为事实来源账务系统、审计

4. 与 RabbitMQ / RocketMQ / Pulsar 对比

维度KafkaRabbitMQRocketMQPulsar
模型分布式日志传统队列(AMQP)日志 + 队列混合日志,存算分离
单机吞吐百万级 msg/s万级十万级十万级
消息回放原生支持困难支持支持
延迟消息不原生支持插件支持原生支持原生支持
路由灵活性弱(按 topic/分区)强(exchange 路由)中(tag 过滤)中
生态最成熟(Connect/Streams/Flink)一般国内电商场景多较新
运维复杂度中(KRaft 后降低)低中高(BookKeeper)

选型建议:

  • 大数据量、日志流、流处理、需要回放 → Kafka
  • 业务消息量不大、需要复杂路由/延迟队列/优先级 → RabbitMQ
  • 电商交易场景、需要事务消息与定时消息 → RocketMQ
  • 多租户、跨地域复制、存算分离诉求强烈 → Pulsar
ℹ️Kafka 不擅长什么

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
🎯练习
  1. 列出你当前系统中三个可以用消息系统解耦的调用链,并说明用 Kafka 还是 RabbitMQ 更合适、为什么。
  2. 思考:Kafka "消费后不删除消息"的设计,对实现"新上线一个数据分析服务、需要重算过去 7 天数据"这类需求有什么帮助?