这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。
Broker 每秒要处理几十万请求,却用一套非常经典的 Reactor 线程模型做到高吞吐低延迟。理解这套模型,num.network.threads、num.io.threads 这些参数才调得有依据。
13.1 线程模型总览
┌────────────────┐
客户端请求 -----> │ Acceptor │ 监听端口,接收连接
└───────┬────────┘
│ 轮询分配
┌───────────────┼───────────────┐
v v v
┌──────────┐ ┌──────────┐ ┌──────────┐
│Processor1│ │Processor2│ │Processor3│ 网络线程
│ (NIO) │ │ (NIO) │ │ (NIO) │ num.network.threads
└────┬─────┘ └────┬─────┘ └────┬─────┘
│ │ │
+-------+-------+------+--------+
v v
┌────────────┐ ┌────────────┐
│ 请求队列 │ │ 响应队列 │
└─────┬──────┘ └─────^──────┘
v │
┌─────────────────────────────────┐
│ RequestHandler 线程池 │ IO 线程
│ num.io.threads │ 真正处理请求
│ (KafkaApis -> ReplicaManager) │
└─────────────────────────────────┘
各角色职责:
- Acceptor:接受 TCP 连接,注册到 Processor;
- Processor(网络线程):读写 Socket、解析/序列化协议,不做业务逻辑;
- 请求队列:待处理请求的缓冲,默认
queued.max.requests=500; - RequestHandler(IO 线程):执行
KafkaApis.handleProduceRequest/handleFetchRequest等,读写日志; - Purgatory:存放“暂时无法完成”的请求,稍后重试检查。
这套模型让网络 IO 与请求处理解耦:慢的磁盘操作不会阻塞网络线程收包,网络突发也不会让 IO 线程无上限堆积。
13.2 Purgatory:延迟操作的“候车室”
有两类典型请求需要等待条件满足:
- Produce 请求
acks=all:Leader 写入本地后,必须等 ISR 其他副本 Fetch 追上,才能给客户端返回成功; - Fetch 请求
fetch.min.bytes:消费者要求“攒够 N 字节再返回”,不足时先挂起,等新数据到达。
如果让 IO 线程自旋等待,吞吐会瞬间崩塌。于是有了 Purgatory:
请求到达 -> 条件不满足 -> 封装成 DelayedOperation 放入 Purgatory
-> 线程立即释放
事件发生(副本追上/新数据写入)-> 尝试 complete 该请求
超时(request.timeout.ms 等) -> 按 timeout 完成/失败
purgatory.size 相关指标过大,说明“等待条件”长期无法满足,常见于副本落后、消费者 min.bytes 过大、acks=all 写入压力大等场景。
13.3 副本拉取线程
Follower 同步走的是独立的 ReplicaFetcherThread:
num.replica.fetchers=1 # 每 Broker 的副本拉取线程数,繁忙集群可调大
replica.fetch.max.bytes=1048576
replica.fetch.wait.max.ms=500
replica.fetch.response.max.bytes=10485760
集群写入流量大、出现 UnderReplicatedPartitions 告警时,适当增大 num.replica.fetchers(如 2-4)常能有效提升副本追赶速度。
13.4 关键 Broker 参数
| 参数 | 默认 | 说明 |
|---|---|---|
num.network.threads |
3 | 网络线程数,CPU 密集 |
num.io.threads |
8 | 请求处理线程数,通常给足 |
queued.max.requests |
500 | 请求队列上限,满则拒绝新请求 |
queued.max.bytes |
-1 | 请求队列字节数限制 |
socket.send.buffer.bytes |
102400 | Socket 发送缓冲 |
socket.receive.buffer.bytes |
102400 | Socket 接收缓冲 |
num.replica.fetchers |
1 | 副本拉取线程数 |
message.max.bytes |
1048588 | 单条消息上限 |
replica.lag.time.max.ms |
30000 | ISR 判定超时 |
调优经验:
NetworkProcessorAvgIdlePercent低于 30%:加大num.network.threads;RequestHandlerAvgIdlePercent低于 30%:加大num.io.threads;- 两者都很空闲但延迟高:问题多半在磁盘/页缓存/客户端,而不是线程数。
13.5 一次 Produce 请求的完整旅程
把前面所有章节串起来:
1. Producer Sender 线程把某 Broker 上多个分区的 batch 合并成 ProduceRequest
2. 网络 -> Broker Acceptor -> Processor 解析 -> 请求队列
3. IO 线程 KafkaApis.handleProduceRequest
4. ReplicaManager 追加到各分区 Leader 本地日志
5. acks=all:请求进入 Purgatory,等待 ISR 副本 Fetch 追上
6. Follower 的 ReplicaFetcherThread 拉取并写入本地
7. Leader 推进 HW,DelayedProduce 完成
8. 响应经响应队列 -> Processor -> 网络 -> Producer 回调
消费 Fetch 类似,但多了零拷贝路径与 fetch.min.bytes 等待。
本章小结
- Broker 是 Reactor 模型:Acceptor + Processor(网络线程) + 请求队列 + IO 线程池;
- Purgatory 让 acks=all 与 fetch.min.bytes 这类延迟请求不占用线程;
- 副本同步由 ReplicaFetcherThread 独立完成,落后时可加线程数;
- 线程参数调优要看 idle 指标,而不是盲调;
- 一条消息的成功确认,横跨网络、日志、副本、协调多条链路。
思考题
- 为什么网络线程不直接处理请求,而要经过请求队列交给 IO 线程?
fetch.min.bytes调大后,延迟和吞吐分别怎么变?它与fetch.max.wait.ms如何配合?- 如果 Purgatory 中积压大量 acks=all 的 Produce 请求,你会优先排查什么?