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 继续消费。

分区分配策略

KafkaConsumerpartition_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.msretention.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


参见

阅读更多

Web 安全基础

1. HTML 转义(服务端渲染必须): 2. CSP(Content Security Policy): 3. HttpOnly Cookie:防止 JS 读取会话 Cookie: 4. 前端框架防护: 攻击者在第三方网站构造一个表单,诱导已登录用户提交,浏览器会自动携带目标站的 Cookie。 触发条件: 1. 用户已登录目标网站(Cookie 有效) 2. 目标 API 仅凭 Cookie 识别用户身份 3. 请求来源未验证 1. CSRF Token(推荐): 2. SameSite Cookie: 3. 验证 Origin/Referer 头:

By yellowdog

HTTP 协议深度指南

HTTP(HyperText Transfer Protocol)是 Web 的基础传输协议,基于 TCP/IP,采用请求/响应模型。 相关文档:Web安全基础(/web-an-quan-ji-chu/) FastAPI完全指南(/fastapi-wan-quan-zhi-nan/) Nginx完全指南(/nginx-wan-quan-zhi-nan/) 幂等性:多次执行相同请求,服务器状态结果相同。PUT /users/1 多次执行结果一致;POST /users 每次创建新资源,非幂等。 浏览器直接从本地缓存读取,不向服务器发送请求。 缓存命中时,状

By yellowdog

系统设计基础

SLA 对照表: 选择建议:无状态服务(Web 层、API 层)优先水平扩展;数据库初期垂直扩展,达到瓶颈后考虑分库分表或读写分离。 缓存穿透(查询不存在的 key,每次都打到 DB): 缓存击穿(热点 key 过期,瞬间大量请求打到 DB): 缓存雪崩(大量 key 同时过期,或缓存服务宕机): 令牌桶 Python 实现: Redis 实现分布式限流(滑动窗口): URL 命名规则: Cursor 分页响应格式: 雪花算法结构(64 bit): 定义:分布式系统不能同时满足以下三个特性: 在分布式环境中 P 是必须保证的,所以实际是 CP vs AP

By yellowdog

算法思路与模板

二分查找要求序列有序,每次将搜索范围缩减一半,时间复杂度 O(log n)。 两个指针从两端向中间收缩,常用于有序数组。 滑动窗口维护一个满足条件的区间 left, right,right 不断向右扩张,条件不满足时收缩 left。 滑动窗口通用框架: 1. 确定"子问题":原问题可以分解为哪些规模更小的同类问题 2. 定义 dpi 或 dpij 的含义,要足够清晰 3. 推导状态转移方程 4. 确定初始状态(边界条件) 5. 确定计算顺序(确保依赖的子问题先计算) 每件物品最多选一次。dpj = 容量为 j 时的最大价值,逆序遍历容量防止重复选取。 每

By yellowdog