Skip to content
Gains Summary
Main Navigation 首页 / Home
C++ 编程 / C++ Programming
系统与高性能 / Systems & Performance
Web 开发 / Web Development
人工智能 / Artificial Intelligence
工业软件 / Industrial Software
其他内容 / Other Topics
C++ 编程 / C++系统与性能 / SystemsWeb 开发 / Web人工智能 / AI工业软件 / Industrial

外观

Sidebar Navigation

← Web 开发 / Web Development

中间件 / Middleware

1. Web 中间件全景 / The Web Middleware Landscape

2. 中间件——流量洪峰下的系统保护 / Middleware for Protecting Systems Under Traffic Spikes

cache layer

1. 缓存架构全景 / Cache Architecture Overview

2. Redis 深入——为什么你的缓存总是出问题 / Redis Internals and Cache Failure Modes

3. Redis 数据结构场景应用——什么时候用什么 / Choosing Redis Data Structures for Real Applications

4. 缓存策略——库存变了,缓存怎么处理 / Cache Strategies and Inventory Consistency

5. Redis 缓存三剑客——穿透/击穿/雪崩 + 一致性策略 / Redis Cache Penetration, Breakdown, Avalanche, and Consistency

6. Redis 分布式锁——从 SETNX 到 Redisson / Redis Distributed Locks from SETNX to Redisson

7. Redis Cluster 与 Sentinel 高可用架构 / Redis Cluster and Sentinel High-Availability Architecture

8. Redis 高级特性——Stream / PubSub / Module / LLM 应用

9. Redis 架构深度分析——为什么 Redis 能这么快 / Redis Architecture and the Sources of Its Performance

10. Redis 全景——为什么你的系统需要一个缓存层 / The Redis Landscape and Why Systems Need a Cache Layer

message queue

1. 消息队列全景——为什么你的系统需要一个中间人 / Message Queue Overview and Why Your System Needs a Middleman

2. Kafka 核心——为什么你的消息总是"丢"了 / Kafka Fundamentals and Message Delivery Semantics

3. 消息队列高级——死信队列、延迟消息、消息积压 / Advanced Messaging with Dead Letters, Delays, and Backlogs

4. Kafka Streams 与 Connect —— 让数据自己流动起来 / Kafka Streams and Connect for Streaming Data Pipelines

5. RabbitMQ 深度解析 —— 灵活路由与消息可靠性 / RabbitMQ Deep Dive into Flexible Routing and Reliability

6. 消息队列对比——为什么最终选了 Kafka / Comparing Message Queues and Choosing Kafka

search engine

1. 搜索引擎知识体系 / Search Engine Knowledge System

2. Elasticsearch——为什么 Like 查询总是那么慢 / Elasticsearch for Full-Text Search at Scale

3. Elasticsearch 查询 DSL 深入——为什么你的搜索总是不准 / Elasticsearch Query DSL Deep Dive

4. Elasticsearch 集群规划与运维——为什么你的集群总是"黄" / Elasticsearch Cluster Planning and Operations

5. Meilisearch 与轻量搜索替代方案——当 ES 太重时 / Meilisearch and Lightweight Search Alternatives

infrastructure

1. 基础设施组件全景 / Infrastructure Components Landscape

2. Nginx 与反向代理 / Nginx and Reverse Proxy

3. 服务发现 / Service Discovery

4. 配置中心 / Configuration Center

本页目录

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

这一章,我们从订单丢失这个问题出发,理解 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

推演:如果你是工程师 ​

第一反应:异步化

用户下单 → 返回"下单成功,正在处理..." → 立刻响应
          ↓
    Kafka(消息队列)
          ↓
    消费者(异步处理)
         ├── 扣减库存
         ├── 创建订单
         ├── 发送通知
         └── 更新 BI

这样:用户响应时间从 250ms → 10ms
     系统吞吐量从受限于业务处理 → 受限于消息队列
     10 万请求排队慢慢处理,不会打垮系统
1
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节:分区——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

写入流程 ​

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

分区数 = 最大并行度 ​

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 1
1
2
3
4
5
6
7
8
9

为什么"按用户 ID 哈希"很重要?

如果用户的两个操作(下单→支付)分别去了不同分区,消费顺序可能是:

分区0: 用户A 下单
分区1: 用户A 支付  ← 先消费了支付(下单还没到)
1
2

按用户 ID 哈希后,两个操作一定在同一分区,顺序一定正确。


第3节:消费者组——多实例并行消费的奥秘 ​

场景:只有一个消费者,消息处理不过来了 ​

10 万消息进入 Kafka
一个消费者处理
→ 消费者崩溃 → 10 万消息无人处理
→ 消费者太慢 → 消息堆积 → 系统延迟
1
2
3
4

推演:如果你是工程师 ​

方案:多个消费者实例,同时消费

消费者实例 1 → 消费 Partition 0 + Partition 1
消费者实例 2 → 消费 Partition 2

速度翻倍!但问题是:消费者实例 3 启动后,分区怎么分配?
1
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

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
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

第4节:消息丢失的三个场景 ​

场景一:生产者没收到 ACK ​

发送消息 → Kafka 返回 ACK → 消息持久化成功
       ↑
       如果 Kafka 还没 ACK,生产者就认为发送成功
       → 消息在 PageCache 中 → 服务器断电 → 消息丢失
1
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

场景二:消费者消费了,但 offset 没提交 ​

python
for message in consumer:
    order = parse_order(message.value)
    process_order(order)  # 处理成功
    # 忘记调用 consumer.commit()
    # Consumer 重启后,从上次提交的 offset 重读
    # 消息被重复消费!
1
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

解决方案:手动提交 + 死信队列

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

第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

推演:如果你是工程师 ​

方案 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

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

第6节:Kafka 事务——原子性保障 ​

场景:扣库存 + 创建订单,必须同时成功或同时失败 ​

普通发送:库存扣了,但订单创建失败了 → 数据不一致

Kafka 事务:
BEGIN TRANSACTION
  send Kafka message (扣库存指令)
  INSERT INTO orders (...)
COMMIT TRANSACTION
→ 要么都成功,要么都回滚
1
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 e
1
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

"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

学习状态:🟡 开始学习

最后更新于:

Pager
上一篇1. 消息队列全景——为什么你的系统需要一个中间人 / Message Queue Overview and Why Your System Needs a Middleman
下一篇3. 消息队列高级——死信队列、延迟消息、消息积压 / Advanced Messaging with Dead Letters, Delays, and Backlogs

持续记录,持续成长

Copyright © Tidenflow