Kafka 核心——为什么你的消息总是"丢"了 / Kafka Fundamentals and Message Delivery Semantics
📅 创建时间:2026-05-08 🏷️ 标签:#Kafka #消息队列 #分区 #消费者组 #幂等性 #消息丢失 📚 前置知识:[[00-backend-overview]] [[01-redis-deep]](理解削峰场景) 📚 相关知识:[[04-distributed-system]](分布式一致性) [[11-architecture-patterns]](CQRS)
场景:你的秒杀订单,消失了
┌─────────────────────────────────────────────────────────────┐
│ │
│ 双十一,10 万用户下单。 │
│ 系统显示"下单成功"。 │
│ │
│ 一小时后,运营发现:只有 9.8 万个订单。 │
│ 丢失了 2000 个订单。 │
│ │
│ 用户已经在下单页等了 3 秒(超时),以为失败了 │
│ 实际后端在"处理中"。 │
│ │
│ 你被叫去排查。 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
2
3
4
5
6
7
8
9
10
11
12
13
14
这一章,我们从订单丢失这个问题出发,理解 Kafka 的核心机制:分区、偏移量、幂等性、事务。
第1节:为什么需要消息队列?
问题抽象
还是秒杀场景。10 万用户同时下单,每个订单需要:
1. 校验库存(Redis)
2. 扣减库存(Redis + MySQL)
3. 创建订单(MySQL)
4. 发送通知(短信/推送)
5. 更新统计(BI 系统)
每个步骤平均耗时 50ms → 总耗时 250ms
→ 10 万 QPS × 250ms = 占用大量线程
→ 数据库连接耗尽
→ 系统崩溃1
2
3
4
5
6
7
8
9
10
2
3
4
5
6
7
8
9
10
推演:如果你是工程师
第一反应:异步化
用户下单 → 返回"下单成功,正在处理..." → 立刻响应
↓
Kafka(消息队列)
↓
消费者(异步处理)
├── 扣减库存
├── 创建订单
├── 发送通知
└── 更新 BI
这样:用户响应时间从 250ms → 10ms
系统吞吐量从受限于业务处理 → 受限于消息队列
10 万请求排队慢慢处理,不会打垮系统1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
2
3
4
5
6
7
8
9
10
11
12
13
14
15
Kafka 的核心价值
┌─────────────────────────────────────────────────────────────┐
│ Kafka 在系统中的角色 │
├─────────────────────────────────────────────────────────────┤
│ │
│ ┌────────────┐ │
│ │ 用户下单 │ → 10 万 QPS(瞬间洪峰) │
│ └─────┬──────┘ │
│ ↓ │
│ ┌────────────┐ │
│ │ Kafka │ ← 削峰填谷:洪峰 → 队列 → 消费者慢慢处理 │
│ │ (消息队列) │ │
│ └─────┬──────┘ │
│ ↓ │
│ ┌────────────┐ ┌────────────┐ ┌────────────┐ │
│ │ 库存服务 │ │ 订单服务 │ │ 通知服务 │ │
│ └────────────┘ └────────────┘ └────────────┘ │
│ │
│ 好处: │
│ • 削峰填谷:不怕洪峰 │
│ • 系统解耦:下单服务和库存服务独立演化 │
│ • 异步处理:用户不用等所有步骤完成 │
│ • 消息持久化:系统崩溃后消息不丢 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
第2节:分区——Kafka 为什么这么快
核心概念
┌─────────────────────────────────────────────────────────────┐
│ Kafka 分区原理 │
├─────────────────────────────────────────────────────────────┤
│ │
│ Topic(主题):逻辑概念,相当于"数据库的表" │
│ Partition(分区):物理概念,相当于"数据库的分区" │
│ Replica(副本):备份,一个分区有多个副本 │
│ │
│ Topic: order-created │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Partition 0: [msg1, msg2, msg3, ...] │ │
│ │ Partition 1: [msg4, msg5, msg6, ...] │ │
│ │ Partition 2: [msg7, msg8, msg9, ...] │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 为什么分区能提速? │
│ • 写:顺序写磁盘(追加写,比随机写快 100 倍) │
│ • 读:并行消费(多个分区可以并行消费) │
│ • 扩展:分区数 = 最大并行度 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
写入流程
Producer Kafka Broker
│ │
│ 1. 查找分区 leader │
│─────────────────────────────────────▶│
│◀─────────────────────────────────────│
│ 返回 Partition 0 的 leader 节点 │
│ │
│ 2. 写入 Partition 0 Leader │
│─────────────────────────────────────▶│
│ (追加写磁盘,PageCache) │
│◀─────────────────────────────────────│
│ 3. 返回 ACK(确认写入成功) │
│ │
│ 4. 异步复制到 Follower │
│ ▼ │
│ ┌───────────────┐ │
│ │ Follower 收到 │ │
│ │ 副本后 ACK │ │
│ └───────────────┘ │1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
分区数 = 最大并行度
python
# 订单消息发到哪个分区?
# 方案 1:轮询(订单均匀分布,但相关订单不在同一分区)
for i in range(10):
kafka.send("order-topic", f"order_{i}", f"订单数据") # 轮流发到 P0,P1,P2
# 方案 2:按用户 ID 哈希(同一用户的订单在同一分区,保证顺序)
kafka.send("order-topic", user_id, f"订单数据_{user_id}")
# user_id=1001 → hash(1001) % 3 = 0 → Partition 0
# user_id=2002 → hash(2002) % 3 = 1 → Partition 11
2
3
4
5
6
7
8
9
2
3
4
5
6
7
8
9
为什么"按用户 ID 哈希"很重要?
如果用户的两个操作(下单→支付)分别去了不同分区,消费顺序可能是:
分区0: 用户A 下单
分区1: 用户A 支付 ← 先消费了支付(下单还没到)1
2
2
按用户 ID 哈希后,两个操作一定在同一分区,顺序一定正确。
第3节:消费者组——多实例并行消费的奥秘
场景:只有一个消费者,消息处理不过来了
10 万消息进入 Kafka
一个消费者处理
→ 消费者崩溃 → 10 万消息无人处理
→ 消费者太慢 → 消息堆积 → 系统延迟1
2
3
4
2
3
4
推演:如果你是工程师
方案:多个消费者实例,同时消费
消费者实例 1 → 消费 Partition 0 + Partition 1
消费者实例 2 → 消费 Partition 2
速度翻倍!但问题是:消费者实例 3 启动后,分区怎么分配?1
2
3
4
5
6
2
3
4
5
6
消费者组原理
┌─────────────────────────────────────────────────────────────┐
│ 消费者组 重平衡 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 消费者组 A(3 个实例): │
│ │
│ Partition 0 ──▶ Consumer 1 │
│ Partition 1 ──▶ Consumer 2 │
│ Partition 2 ──▶ Consumer 3 │
│ │
│ 新 Consumer 4 加入: │
│ │
│ ┌─────────────────────────────────────────────────────┐ │
│ │ Rebalance 触发! │ │
│ │ Partition 0 ──▶ Consumer 1 (重新分配) │ │
│ │ Partition 1 ──▶ Consumer 4 (新分配的) │ │
│ │ Partition 2 ──▶ Consumer 3 │ │
│ │ Consumer 2 暂时不消费 │ │
│ └─────────────────────────────────────────────────────┘ │
│ │
│ 问题:Rebalance 期间消息可能重复消费 │
│ Consumer 2 正在处理 msg5 → Rebalance 触发 → msg5 可能被 Consumer 4 重新消费│
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
offset 管理:消息有没有被消费?
┌─────────────────────────────────────────────────────────────┐
│ Consumer Offset 管理 │
├─────────────────────────────────────────────────────────────┤
│ │
│ offset = 消费者读取到哪个位置了 │
│ │
│ Consumer 处理消息的过程: │
│ │
│ Partition 0: [msg0][msg1][msg2][msg3][msg4][msg5]... │
│ │
│ 已提交 offset: 2 │
│ Consumer 读取: msg3 → 处理中... │
│ │
│ offset 提交时机: │
│ • 自动提交(默认):每 5 秒提交一次 │
│ → Consumer 处理 msg3 中途崩溃 → 重启后从 offset=2 重读│
│ → msg3 被重复消费(幂等性问题!) │
│ │
│ • 手动提交:处理完 msg3 后再提交 offset=3 │
│ → Consumer 处理 msg3 中途崩溃 → 重启后从 offset=3 重读 │
│ → msg3 不会重复消费(但可能丢失) │
│ │
│ 最佳实践:处理成功后再提交 offset │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
python
from kafka import KafkaConsumer
consumer = KafkaConsumer(
'order-topic',
bootstrap_servers=['kafka:9092'],
group_id='order-processor', # 消费者组名
enable_auto_commit=False, # 关闭自动提交
auto_offset_reset='earliest' # 无 offset 时从最早开始
)
for message in consumer:
try:
order = parse_order(message.value)
# 业务处理
process_order(order)
# 处理成功后再提交 offset
consumer.commit()
except Exception as e:
# 处理失败,不提交 offset,下次重启后重新消费
print(f"处理失败: {e}, offset={message.offset}")1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
第4节:消息丢失的三个场景
场景一:生产者没收到 ACK
发送消息 → Kafka 返回 ACK → 消息持久化成功
↑
如果 Kafka 还没 ACK,生产者就认为发送成功
→ 消息在 PageCache 中 → 服务器断电 → 消息丢失1
2
3
4
2
3
4
解决方案:确认机制
python
# acks=0:发送即成功(最快,最容易丢)
producer.send("topic", value=data) # fire and forget
# acks=1:Leader 写入成功即返回(折中)
producer.send("topic", value=data, acks=1)
# acks=all:Leader + 所有 ISR 副本写入成功才返回(最安全)
producer.send("topic", value=data, acks="all")
# 重试机制
producer.send("topic", value=data, acks="all", retries=3)1
2
3
4
5
6
7
8
9
10
11
2
3
4
5
6
7
8
9
10
11
场景二:消费者消费了,但 offset 没提交
python
for message in consumer:
order = parse_order(message.value)
process_order(order) # 处理成功
# 忘记调用 consumer.commit()
# Consumer 重启后,从上次提交的 offset 重读
# 消息被重复消费!1
2
3
4
5
6
2
3
4
5
6
解决方案:手动提交 offset + 幂等处理
场景三:消费者消费了,但处理失败了
python
for message in consumer:
try:
order = parse_order(message.value)
process_order(order) # 这里抛异常了
consumer.commit() # 这行没执行
except:
pass # 异常被吞了
# Consumer 重启后,从上次提交的 offset 重读
# 消息丢失了!1
2
3
4
5
6
7
8
9
2
3
4
5
6
7
8
9
解决方案:手动提交 + 死信队列
python
for message in consumer:
try:
order = parse_order(message.value)
process_order(order)
consumer.commit()
except Exception as e:
# 发送到死信队列(人工处理)
send_to_dlq("order-topic", message.value, str(e))
consumer.commit() # 确认消费,防止无限重试1
2
3
4
5
6
7
8
9
2
3
4
5
6
7
8
9
第5节:幂等性——如何保证消息不被重复消费
场景:用户手抖,多点了一次下单
用户快速点了两次"立即购买"
消息 1: order_id=A, user_id=1, product_id=100, amount=1
消息 2: order_id=B, user_id=1, product_id=100, amount=1
两个消息都被消费了
→ 创建了两个订单
→ 用户被扣了两次钱
→ 投诉!1
2
3
4
5
6
7
8
9
2
3
4
5
6
7
8
9
推演:如果你是工程师
方案 1:数据库唯一索引(最常用)
CREATE TABLE orders (
id BIGINT PRIMARY KEY,
idempotency_key VARCHAR(64) UNIQUE, -- 幂等键
user_id BIGINT,
...
);
INSERT INTO orders (...) VALUES (...) ON CONFLICT (idempotency_key) DO NOTHING;
-- 第二个请求会冲突,直接忽略
方案 2:Redis 记录已处理的 key
if redis.exists(f"processed:{idempotency_key}"):
return # 重复消费,直接返回
process_order(order)
redis.setex(f"processed:{idempotency_key}", 86400, "1") # 24 小时内不重复处理
方案 3:Kafka 幂等生产者
producer = KafkaProducer(
"order-topic",
enable_idempotence=True # Kafka 自动保证幂等
)
# Kafka 内部会记录每个 Producer 发送的消息
# 重复发送的消息,Kafka 会丢弃1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
Kafka 幂等生产者原理
┌─────────────────────────────────────────────────────────────┐
│ Kafka 幂等生产者原理 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 每个 Producer 有一个唯一的 PID(Producer ID) │
│ 每次发送消息,带上 PID + Sequence Number │
│ │
│ Kafka Broker 记录: │
│ PID=101, Seq=5 ← 最后处理的消息序号 │
│ │
│ Producer 发送消息: │
│ PID=101, Seq=5 → Broker 发现 seq=5 已处理过 → 丢弃 │
│ PID=101, Seq=6 → Broker 处理 Seq=6 │
│ PID=102, Seq=1 → 新 Producer,从 1 开始 │
│ │
│ 限制:只能保证单个 Producer 的幂等 │
│ 如果 Producer 重启,PID 变了,幂等性不跨 Producer │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
第6节:Kafka 事务——原子性保障
场景:扣库存 + 创建订单,必须同时成功或同时失败
普通发送:库存扣了,但订单创建失败了 → 数据不一致
Kafka 事务:
BEGIN TRANSACTION
send Kafka message (扣库存指令)
INSERT INTO orders (...)
COMMIT TRANSACTION
→ 要么都成功,要么都回滚1
2
3
4
5
6
7
8
2
3
4
5
6
7
8
python
from kafka import KafkaProducer
from kafka.errors import KafkaError
producer = KafkaProducer(bootstrap_servers=['kafka:9092'])
# 开启事务
producer.init_transactions()
try:
producer.begin_transaction()
# 发送扣库存消息
producer.send("inventory-topic", value={"product_id": 100, "delta": -1})
producer.send("inventory-topic", value={"product_id": 101, "delta": -1})
# 提交事务:所有消息要么全部发送成功,要么全部不发送
producer.commit_transaction()
except KafkaError as e:
# 任意一个失败,全部回滚
producer.abort_transaction()
raise e1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
升华:Kafka 的核心设计哲学
┌─────────────────────────────────────────────────────────────┐
│ Kafka 的三个核心权衡 │
├─────────────────────────────────────────────────────────────┤
│ │
│ 1. 顺序写 vs 随机读 │
│ Kafka 用追加写(顺序写磁盘)→ 写入极快 │
│ 读取时用 PageCache(内存)→ 读取也不慢 │
│ │
│ 2. at-least-once vs at-most-once vs exactly-once │
│ at-least-once:可能重复消费,不会丢(常用) │
│ at-most-once:可能丢消息,不会重复 │
│ exactly-once:Kafka 内部可保证,跨系统需事务(最贵) │
│ │
│ 3. 持久化 vs 性能 │
│ Kafka 将消息持久化到磁盘(不怕宕机) │
│ 但通过 PageCache + 顺序写 + 零拷贝, │
│ 性能依然极高(比随机写内存还快) │
│ │
│ 一句话总结: │
│ Kafka 是用磁盘实现了内存的性能,同时保证了持久化。 │
│ │
└─────────────────────────────────────────────────────────────┘1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
"AI 可查 vs 必须理解"清单
AI 可查:
✅ Kafka 命令行工具(kafka-topics.sh / kafka-console-producer.sh)
✅ Python/Java/Go 的 Kafka 客户端 API
✅ 分区分配策略的具体参数
✅ 副本 ISR 的配置参数
必须理解:
🔴 分区数和消费者数的关系(不是越多越好)
🔴 offset 提交的时机和幂等性的关系
🔴 at-least-once / at-most-once / exactly-once 的区别
🔴 消息丢失的三个场景(生产者/Broker/消费者)
🔴 为什么按用户 ID 哈希分区能保证同一用户的顺序1
2
3
4
5
6
7
8
9
10
11
12
2
3
4
5
6
7
8
9
10
11
12
学习状态:🟡 开始学习