KafkaNotes

第 13 章:请求处理模型:网络线程与 Purgatory

zjc 于 2026-01-13 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Broker 每秒要处理几十万请求,却用一套非常经典的 Reactor 线程模型做到高吞吐低延迟。理解这套模型,num.network.threadsnum.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)  │
        └─────────────────────────────────┘

各角色职责:

这套模型让网络 IO 与请求处理解耦:慢的磁盘操作不会阻塞网络线程收包,网络突发也不会让 IO 线程无上限堆积。

13.2 Purgatory:延迟操作的“候车室”

有两类典型请求需要等待条件满足:

  1. Produce 请求 acks=all:Leader 写入本地后,必须等 ISR 其他副本 Fetch 追上,才能给客户端返回成功;
  2. 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 判定超时

调优经验:

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 等待。

本章小结

思考题

  1. 为什么网络线程不直接处理请求,而要经过请求队列交给 IO 线程?
  2. fetch.min.bytes 调大后,延迟和吞吐分别怎么变?它与 fetch.max.wait.ms 如何配合?
  3. 如果 Purgatory 中积压大量 acks=all 的 Produce 请求,你会优先排查什么?