Kafka 完全指南
Apache Kafka 是一个分布式流处理平台,以高吞吐量、持久化存储和消息回放著称。 KafkaProducer 参数: producer.send 参数: KafkaConsumer 参数: acks 控制生产者认为消息发送成功所需的 Broker 确认数量: 生产环境金融/订单类消息应使用 acks="all" + retries=3。 注意:auto_offset_reset 只在该 Consumer Group 从未提交过 offset 时生效(新 group 或 offset 已过期被删除)。已有提交记录的情况下,从上次提交的 offset
官方文档: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(同步)
安装
pip install kafka-python
生产者(KafkaProducer)
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)
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(异步)
安装
pip install aiokafka
异步生产者
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())
异步消费者
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 集成
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,导致消息丢失(至多一次语义)。手动提交结合业务幂等是实现至少一次语义的标准做法。
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。