> ## Content Index
> Fetch the complete content index at: https://blog.vercanti.com/llms.txt
> Use this file to discover other available public pages before exploring further.

# Kafka 完全指南
- URL: https://blog.vercanti.com/kafka-wan-quan-zhi-nan/
- Published: 2026-08-28T14:35:41.000Z
- Updated: 2026-08-28T14:59:22.000Z
- Description: Apache Kafka 是一个分布式流处理平台，以高吞吐量、持久化存储和消息回放著称。 KafkaProducer 参数： producer.send 参数： KafkaConsumer 参数： acks 控制生产者认为消息发送成功所需的 Broker 确认数量： 生产环境金融/订单类消息应使用 acks="all" + retries=3。 注意：auto_offset_reset 只在该 Consumer Group 从未提交过 offset 时生效（新 group 或 offset 已过期被删除）。已有提交记录的情况下，从上次提交的 offset
- Author: yellowdog
- Tags: 消息队列

> 官方文档：<https://kafka.apache.org/documentation/>  
> 适用版本：Kafka 3.x（2026-05-07 整理）

Apache Kafka 是一个分布式流处理平台，以高吞吐量、持久化存储和消息回放著称。

---

## 核心概念

### 组件关系

```
Producer
    |
    | 发布消息到 Topic（可指定 partition key）
    v
Broker 集群（多个 Broker 组成）
    |
    | Topic 拆分为多个 Partition，分布在不同 Broker 上
    | 每个 Partition 有一个 Leader 和多个 Follower（Replication）
    v
Partition（有序、不可变的消息日志）
    |
    | Consumer Group 中每个 Partition 只分配给一个 Consumer
    v
Consumer（Consumer Group 内的成员）

```

- **Topic**：消息的逻辑分类，类似数据库中的表名。
- **Partition**：Topic 的物理分片，每个 Partition 是一个有序的追加日志。分区数决定并行消费能力。
- **Offset**：消息在 Partition 内的唯一序号，从 0 递增。Consumer 通过 offset 追踪消费进度。
- **Consumer Group**：一组 Consumer 共享同一个 `group_id`，每条消息只被组内一个 Consumer 处理。不同 Consumer Group 之间互相独立，均可消费全量消息。
- **Broker**：Kafka 集群中的单个节点，负责存储和转发消息。
- **Replication**：Partition 副本机制，Leader 处理读写，Follower 同步数据。Leader 故障时自动选举新 Leader。

### 与 RabbitMQ 的设计差异

| 维度   | Kafka                                       | RabbitMQ               |
| ---- | ------------------------------------------- | ---------------------- |
| 存储模型 | 顺序追加日志（Log），消息保留至配置的时间/大小上限                 | 队列，消息被 ack 后删除         |
| 消费模型 | Consumer 主动 Pull，自行管理 offset                | Broker Push 给 Consumer |
| 消息回放 | 支持，重置 offset 可重新消费历史消息                      | 不支持                    |
| 扩展方式 | 增加 Partition 数量，Consumer Group 自动 Rebalance | 增加队列和消费者               |
| 吞吐量  | 极高（百万级 QPS），适合流数据                           | 高（万级 QPS），适合任务队列       |
| 消息路由 | 简单，靠 topic + key 哈希分区                       | 复杂，支持多种 Exchange 类型    |

---

## Python 客户端：kafka-python（同步）

### 安装

```bash
pip install kafka-python

```

### 生产者（KafkaProducer）

```python
from kafka import KafkaProducer
import json

producer = KafkaProducer(
    bootstrap_servers=["localhost:9092"],
    value_serializer=lambda v: json.dumps(v).encode("utf-8"),
    key_serializer=lambda k: k.encode("utf-8") if k else None,
    acks="all",
    retries=3,
)

# 发送消息
future = producer.send(
    topic="orders",
    value={"order_id": 123, "amount": 99.9},
    key="user_456",
)
record_metadata = future.get(timeout=10)  # 阻塞等待确认
print(f"partition={record_metadata.partition}, offset={record_metadata.offset}")

producer.flush()   # 确保所有缓冲消息已发送
producer.close()

```

`KafkaProducer` 参数：

| 参数                   | 类型       | 默认值      | 说明                                   |
| -------------------- | -------- | -------- | ------------------------------------ |
| bootstrap\_servers   | list/str | 必填       | Broker 地址列表，用于初始连接，之后自动发现集群          |
| value\_serializer    | callable | None     | 消息 value 序列化函数，接收原始值，返回 bytes        |
| key\_serializer      | callable | None     | 消息 key 序列化函数                         |
| acks                 | int/str  | 1        | 消息确认级别：0/1/"all"，见下方说明               |
| retries              | int      | 0        | 发送失败重试次数                             |
| batch\_size          | int      | 16384    | 批量发送的最大字节数（16KB），提升吞吐量               |
| linger\_ms           | int      | 0        | 消息在缓冲区等待批量发送的最长时间（毫秒）。增大此值可提高吞吐但增加延迟 |
| buffer\_memory       | int      | 33554432 | 生产者缓冲区总大小（32MB）                      |
| compression\_type    | str      | None     | 压缩类型："gzip" / "snappy" / "lz4"       |
| max\_block\_ms       | int      | 60000    | send() 阻塞等待缓冲区可用的最长时间                |
| request\_timeout\_ms | int      | 30000    | 请求超时时间（毫秒）                           |

`producer.send` 参数：

| 参数            | 类型   | 默认值  | 说明                                          |
| ------------- | ---- | ---- | ------------------------------------------- |
| topic         | str  | 必填   | 目标 Topic 名称                                 |
| value         | any  | None | 消息体，经 value\_serializer 序列化后发送              |
| key           | any  | None | 消息键，用于分区路由（相同 key 发往同一 partition）           |
| partition     | int  | None | 强制指定分区，覆盖 key 路由策略                          |
| timestamp\_ms | int  | None | 消息时间戳，默认使用当前时间                              |
| headers       | list | None | 消息头，格式为 \[(key\_bytes, value\_bytes), ...\] |

### 消费者（KafkaConsumer）

```python
from kafka import KafkaConsumer
import json

consumer = KafkaConsumer(
    "orders",
    bootstrap_servers=["localhost:9092"],
    group_id="order_processors",
    auto_offset_reset="earliest",
    enable_auto_commit=False,
    value_deserializer=lambda v: json.loads(v.decode("utf-8")),
)

for message in consumer:
    print(f"topic={message.topic}, partition={message.partition}, "
          f"offset={message.offset}, value={message.value}")
    consumer.commit()  # 手动提交 offset

```

`KafkaConsumer` 参数：

| 参数                         | 类型       | 默认值      | 说明                                             |
| -------------------------- | -------- | -------- | ---------------------------------------------- |
| \*topics                   | str      | 可选       | 订阅的 Topic 名称，也可用 subscribe() 动态设置              |
| bootstrap\_servers         | list/str | 必填       | Broker 地址列表                                    |
| group\_id                  | str      | None     | 消费者组 ID，相同 group\_id 的消费者共享分区                  |
| auto\_offset\_reset        | str      | "latest" | 无已提交 offset 时的起始位置："earliest" 从最早，"latest" 从最新 |
| enable\_auto\_commit       | bool     | True     | 是否自动提交 offset，生产环境建议 False 手动控制                |
| auto\_commit\_interval\_ms | int      | 5000     | 自动提交 offset 的时间间隔（毫秒）                          |
| value\_deserializer        | callable | None     | 消息 value 反序列化函数                                |
| key\_deserializer          | callable | None     | 消息 key 反序列化函数                                  |
| session\_timeout\_ms       | int      | 10000    | 消费者心跳超时，超时则触发 Rebalance                        |
| heartbeat\_interval\_ms    | int      | 3000     | 心跳发送间隔，应小于 session\_timeout\_ms 的 1/3          |
| max\_poll\_records         | int      | 500      | 单次 poll() 返回的最大消息数                             |
| max\_poll\_interval\_ms    | int      | 300000   | 两次 poll() 之间允许的最大间隔，超时视为消费者失活                  |

---

## Python 客户端：aiokafka（异步）

### 安装

```bash
pip install aiokafka

```

### 异步生产者

```python
from aiokafka import AIOKafkaProducer
import json
import asyncio

async def produce():
    producer = AIOKafkaProducer(
        bootstrap_servers="localhost:9092",
        value_serializer=lambda v: json.dumps(v).encode(),
    )
    await producer.start()
    try:
        await producer.send_and_wait(
            "orders",
            value={"order_id": 1},
            key=b"user_1",
        )
    finally:
        await producer.stop()

asyncio.run(produce())

```

### 异步消费者

```python
from aiokafka import AIOKafkaConsumer
import json
import asyncio

async def consume():
    consumer = AIOKafkaConsumer(
        "orders",
        bootstrap_servers="localhost:9092",
        group_id="order_processors",
        auto_offset_reset="earliest",
        enable_auto_commit=False,
        value_deserializer=lambda v: json.loads(v.decode()),
    )
    await consumer.start()
    try:
        async for message in consumer:
            print(f"offset={message.offset}, value={message.value}")
            await consumer.commit()
    finally:
        await consumer.stop()

asyncio.run(consume())

```

### 与 FastAPI 集成

```python
from contextlib import asynccontextmanager
from aiokafka import AIOKafkaProducer, AIOKafkaConsumer
from fastapi import FastAPI
import asyncio
import json

producer: AIOKafkaProducer = None

@asynccontextmanager
async def lifespan(app: FastAPI):
    global producer
    producer = AIOKafkaProducer(
        bootstrap_servers="localhost:9092",
        value_serializer=lambda v: json.dumps(v).encode(),
    )
    await producer.start()

    # 后台消费任务
    consumer_task = asyncio.create_task(start_consumer())

    yield

    await producer.stop()
    consumer_task.cancel()

async def start_consumer():
    consumer = AIOKafkaConsumer(
        "orders",
        bootstrap_servers="localhost:9092",
        group_id="api_consumers",
    )
    await consumer.start()
    try:
        async for message in consumer:
            await handle_message(message)
    finally:
        await consumer.stop()

app = FastAPI(lifespan=lifespan)

@app.post("/orders")
async def create_order(order: dict):
    await producer.send_and_wait("orders", value=order)
    return {"status": "queued"}

```

---

## 核心配置

### 生产者 acks 配置

`acks` 控制生产者认为消息发送成功所需的 Broker 确认数量：

| acks 值      | 含义                       | 数据安全性                            | 性能 |
| ----------- | ------------------------ | -------------------------------- | -- |
| 0           | 不等待任何确认，Fire and Forget  | 最低，消息可能丢失                        | 最高 |
| 1           | 等待 Partition Leader 写入确认 | 中等，Leader 故障时未同步到 Follower 的消息丢失 | 较高 |
| "all" / \-1 | 等待所有 ISR（同步副本集合）确认       | 最高，只要有一个副本存活消息不丢失                | 较低 |

生产环境金融/订单类消息应使用 `acks="all"` \+ `retries=3`。

### 消费者 auto\_offset\_reset

| 值          | 含义                | 适用场景             |
| ---------- | ----------------- | ---------------- |
| "earliest" | 从该分区最早可用的消息开始消费   | 需要处理所有历史消息，如数据迁移 |
| "latest"   | 从消费者启动后新产生的消息开始消费 | 只关心实时数据，不需要历史消息  |

注意：`auto_offset_reset` 只在该 Consumer Group 从未提交过 offset 时生效（新 group 或 offset 已过期被删除）。已有提交记录的情况下，从上次提交的 offset 继续消费。

### 分区分配策略

`KafkaConsumer` 的 `partition_assignment_strategy` 参数（默认 `RangeAssignor`）：

| 策略                 | 分配规则                         | 特点                   |
| ------------------ | ---------------------------- | -------------------- |
| RangeAssignor      | 按分区范围分配，同一 Topic 的分区连续分配给消费者 | 分配不均衡（分区数不能被消费者数整除时） |
| RoundRobinAssignor | 轮询分配所有 Topic 的所有分区           | 分配更均衡                |
| StickyAssignor     | 尽量保持原有分配，仅调整必要的分区            | Rebalance 开销最小，推荐使用  |

---

## 消费者组

### 分区分配规则

- 一个 Consumer Group 内，每个 Partition 只能分配给一个 Consumer。
- 若 Consumer 数量多于 Partition 数量，多余的 Consumer 闲置，不消费任何分区。
- 若 Consumer 数量少于 Partition 数量，部分 Consumer 消费多个 Partition。
- 最大并行度 = Partition 数量，因此 Partition 数量决定了消费的水平扩展上限。

### Rebalance 触发条件

Rebalance 是 Consumer Group 内部重新分配 Partition 的过程，触发期间所有消费者暂停消费：

- Consumer 加入 Group（新消费者启动）
- Consumer 离开 Group（消费者正常关闭）
- Consumer 心跳超时（`session_timeout_ms` 内未发送心跳）
- Consumer 两次 `poll()` 间隔超过 `max_poll_interval_ms`（消息处理过慢）
- Topic 的 Partition 数量变化

### 心跳与超时配置

| 参数                      | 默认值    | 说明                                   |
| ----------------------- | ------ | ------------------------------------ |
| session\_timeout\_ms    | 10000  | Consumer 被认为存活的超时时间，Broker 端判断       |
| heartbeat\_interval\_ms | 3000   | 心跳发送间隔，应为 session\_timeout\_ms 的 1/3 |
| max\_poll\_interval\_ms | 300000 | 两次 poll() 之间的最大允许间隔，超时触发 Rebalance   |

若单条消息处理时间较长（如调用外部接口），应增大 `max_poll_interval_ms` 或减少 `max_poll_records`，避免因处理超时导致无意义的 Rebalance。

---

## 踩坑与注意事项

**消息乱序**：Kafka 只保证单个 Partition 内的消息有序，跨 Partition 不保证顺序。若业务需要有序处理（如同一用户的操作日志），必须将相关消息发到同一 Partition（通过相同的 `key`，Kafka 对 key 做哈希路由）。增加 Partition 数量不会破坏同 key 消息的相对顺序。

**Rebalance 风暴**：消费者处理消息耗时过长，导致超过 `max_poll_interval_ms` 被踢出 Group，触发 Rebalance；Rebalance 期间消费暂停，消息积压；积压后消费者再次超时，形成恶性循环。解决方案：减小 `max_poll_records`、增大 `max_poll_interval_ms`、或将耗时操作异步化。

**消息幂等处理**：消费者在处理消息后、提交 offset 前崩溃，重启后会重新消费同一批消息（至少一次语义）。业务逻辑必须实现幂等，常用方案：数据库唯一约束、Redis 去重（以消息 offset 或业务 ID 为 key）。

**offset 提交时机**：使用 `enable_auto_commit=False` 手动提交时，应在消息处理完成后再调用 `commit()`，而不是处理前提交。处理前提交可能导致消息处理失败后该消息被跳过（至多一次语义）。

**Consumer Group ID 冲突**：多个不同业务的消费者使用了相同的 `group_id` 会导致分区被意外分配，部分消费者收不到消息。应为每个独立的消费业务使用唯一的 `group_id`，并在命名上体现业务含义（如 `payment-service-consumer`）。

---

## 最佳实践

**Partition 数量与消费者数量对齐，不要超过消费者数量**：多余的消费者永远处于闲置状态，既浪费资源又不增加并发，反而在 Rebalance 时增加代价。初始 Partition 数量可以是预期消费者数量的 2 倍，以预留扩展空间。

**手动提交 offset（`enable_auto_commit=False`），消息处理成功后再提交**：自动提交可能在消息处理失败前就已提交 offset，导致消息丢失（至多一次语义）。手动提交结合业务幂等是实现至少一次语义的标准做法。

```python
consumer = KafkaConsumer(
    'orders',
    group_id='order-processor',
    enable_auto_commit=False,
    bootstrap_servers=['localhost:9092'],
)
for msg in consumer:
    process(msg)            # 先处理
    consumer.commit()       # 后提交

```

**为消息 key 选择业务含义字段，保证同类消息路由到同一 Partition**：订单消息以 `order_id` 为 key、用户操作以 `user_id` 为 key，确保同一业务实体的消息被同一消费者顺序处理，同时保持整体高吞吐。

**为 Topic 设置 `retention.ms` 和 `retention.bytes`，防止磁盘耗尽**：默认保留 7 天；生产环境应根据业务需求和磁盘容量明确设置，并配置监控告警。

**Consumer 的 `max_poll_interval_ms` 要大于单条消息的最大处理时间**：否则处理稍慢就触发 Rebalance，形成恶性循环。如有慢消息，应将耗时操作（HTTP 请求、数据库写入）异步化，poll 线程只做轻量协调。

---

## 常见陷阱

### 陷阱：消费者数量超过 Partition 数量，部分消费者永远闲置

**现象：** 增加了更多消费者实例，但消费速度没有提升，日志显示部分消费者没有分配到任何 Partition。

**原因：** Kafka 的分配规则是一个 Partition 最多分配给 Consumer Group 内的一个 Consumer；Consumer 数量超过 Partition 数量时，多余的消费者无法参与消费。

**解决：** 增加 Topic 的 Partition 数量，或减少消费者实例数量使其不超过 Partition 数。注意增加 Partition 数量会触发一次 Rebalance。

### 陷阱：处理耗时超过 `max_poll_interval_ms` 引发 Rebalance 风暴

**现象：** 消费者日志反复出现 "Kicked out of the group, triggered rebalance"，消息积压持续增长，消费者一直在重平衡而无法正常消费。

**原因：** 消费者在两次 `poll()` 之间的处理时间超过 `max_poll_interval_ms`（默认 5 分钟），Broker 认为该消费者已死亡，将其踢出 Group 并触发 Rebalance。

**解决：** 减少 `max_poll_records`（每次 poll 的消息数）、增大 `max_poll_interval_ms`，或将耗时操作（如外部接口调用）移到异步线程池，poll 线程快速返回。

### 陷阱：相同 `group_id` 被多个不同业务共用

**现象：** 某些消息被错误的消费者处理，或某个消费者完全收不到预期的消息。

**原因：** 同一 `group_id` 的所有消费者共享分区，Kafka 把分区分配给 Group 内的消费者，不同业务共用 ID 导致分区分配错乱。

**解决：** 每个独立业务使用唯一的 `group_id`，命名规范如 `{服务名}-{topic名}-consumer`。

---

## 参见

- [RabbitMQ完全指南](https://blog.vercanti.com/rabbitmq-wan-quan-zhi-nan/)
- [Redis完全指南](https://blog.vercanti.com/redis-wan-quan-zhi-nan/)
- [FastAPI完全指南](https://blog.vercanti.com/fastapi-wan-quan-zhi-nan/)