第24讲 | 请求处理全链路:网络线程与 IO 线程模型
导读:一个 FETCH 请求打进 broker 后,要穿过 Acceptor、Processor(网络线程)、请求队列、KafkaRequestHandler(IO 线程)才能碰到日志文件。理解这套"Reactor + 工作线程池"的分层,你才能解释"broker 线程数怎么配、请求为什么会排队、purgeable/processing 时间指标看什么"。
本讲目标
一、总览:一张图看懂 broker 内部线程布局
#mermaid-svg-ZToJoEqriTww0Udg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-ZToJoEqriTww0Udg .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-ZToJoEqriTww0Udg .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-ZToJoEqriTww0Udg .error-icon{fill:#552222;}#mermaid-svg-ZToJoEqriTww0Udg .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-ZToJoEqriTww0Udg .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-ZToJoEqriTww0Udg .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-ZToJoEqriTww0Udg .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-ZToJoEqriTww0Udg .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-ZToJoEqriTww0Udg .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-ZToJoEqriTww0Udg .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-ZToJoEqriTww0Udg .marker{fill:#333333;stroke:#333333;}#mermaid-svg-ZToJoEqriTww0Udg .marker.cross{stroke:#333333;}#mermaid-svg-ZToJoEqriTww0Udg svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-ZToJoEqriTww0Udg p{margin:0;}#mermaid-svg-ZToJoEqriTww0Udg .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-ZToJoEqriTww0Udg .cluster-label text{fill:#333;}#mermaid-svg-ZToJoEqriTww0Udg .cluster-label span{color:#333;}#mermaid-svg-ZToJoEqriTww0Udg .cluster-label span p{background-color:transparent;}#mermaid-svg-ZToJoEqriTww0Udg .label text,#mermaid-svg-ZToJoEqriTww0Udg span{fill:#333;color:#333;}#mermaid-svg-ZToJoEqriTww0Udg .node rect,#mermaid-svg-ZToJoEqriTww0Udg .node circle,#mermaid-svg-ZToJoEqriTww0Udg .node ellipse,#mermaid-svg-ZToJoEqriTww0Udg .node polygon,#mermaid-svg-ZToJoEqriTww0Udg .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-ZToJoEqriTww0Udg .rough-node .label text,#mermaid-svg-ZToJoEqriTww0Udg .node .label text,#mermaid-svg-ZToJoEqriTww0Udg .image-shape .label,#mermaid-svg-ZToJoEqriTww0Udg .icon-shape .label{text-anchor:middle;}#mermaid-svg-ZToJoEqriTww0Udg .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-ZToJoEqriTww0Udg .rough-node .label,#mermaid-svg-ZToJoEqriTww0Udg .node .label,#mermaid-svg-ZToJoEqriTww0Udg .image-shape .label,#mermaid-svg-ZToJoEqriTww0Udg .icon-shape .label{text-align:center;}#mermaid-svg-ZToJoEqriTww0Udg .node.clickable{cursor:pointer;}#mermaid-svg-ZToJoEqriTww0Udg .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-ZToJoEqriTww0Udg .arrowheadPath{fill:#333333;}#mermaid-svg-ZToJoEqriTww0Udg .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-ZToJoEqriTww0Udg .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-ZToJoEqriTww0Udg .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-ZToJoEqriTww0Udg .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-ZToJoEqriTww0Udg .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-ZToJoEqriTww0Udg .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-ZToJoEqriTww0Udg .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-ZToJoEqriTww0Udg .cluster text{fill:#333;}#mermaid-svg-ZToJoEqriTww0Udg .cluster span{color:#333;}#mermaid-svg-ZToJoEqriTww0Udg div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-ZToJoEqriTww0Udg .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-ZToJoEqriTww0Udg rect.text{fill:none;stroke-width:0;}#mermaid-svg-ZToJoEqriTww0Udg .icon-shape,#mermaid-svg-ZToJoEqriTww0Udg .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-ZToJoEqriTww0Udg .icon-shape p,#mermaid-svg-ZToJoEqriTww0Udg .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-ZToJoEqriTww0Udg .icon-shape .label rect,#mermaid-svg-ZToJoEqriTww0Udg .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-ZToJoEqriTww0Udg .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-ZToJoEqriTww0Udg .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-ZToJoEqriTww0Udg :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
轮转分发 NIO 连接
take 请求
需要等待时
条件满足后
写回 Socket
客户端连接
Acceptor 线程每 listener 1 个监听端口 接受连接
Processor-0网络线程
Processor-1网络线程
Processor-N网络线程
RequestChannel.requestQueue共享请求队列 默认500
KafkaRequestHandler-0IO线程池
H2 Handler-1
H2
Handler-M
KafkaApis路由到 produce/fetch/… 处理器
日志读写/副本拉取/元数据操作
Purgatory 延迟操作第27讲
Response 队列
客户端
三种线程的分工哲学:收发字节的(网络线程)、执行业务逻辑的(IO 线程)、等待外部条件的(purgatory)彻底分离,各自独立扩缩、互不阻塞。
二、Acceptor 与 Processor:网络层
2.1 Acceptor
- 每个 listener.name 一个 Acceptor 线程,阻塞在 ServerSocketChannel.accept();
- 新连接建立后,以 round-robin 方式交给 Processor;
- 自己不做任何读写——单线程也扛得住,因为 accept 本身极轻。
2.2 Processor(网络线程,num.network.threads)
Processor 是 Reactor 模式里的 SubReactor,基于 Java NIO:
为什么网络线程不顺便把活干了? 因为请求处理时间方差极大:一次元数据请求微秒级,一次大 FETCH 可能要等 purgatory 数百毫秒。若在同一批线程里做,一个慢请求就会拖住所有连接的收发,连心跳和元数据请求都会超时——进而引发"客户端集体 rebalance/断连"的雪崩。分层之后,慢请求只会堆积在请求队列里,网络层始终轻快。
2.3 关键参数
| num.network.threads | 3 | 每 listener 的 Processor 数 |
| queued.max.requests | 500 | 共享请求队列上限,超出后网络线程停止读取(背压!) |
| queued.request.max.bytes(broker 请求级还有数据量维度) | – | 3.x 对请求体大小另有保护 |
| connections.max.idle.ms | 10 分钟 | 空闲连接回收 |
| listener.security.protocol.map | – | listener 与协议映射 |
注意 queued.max.requests 的背压语义:队列满时 Processor 不再从 socket 读数据,TCP 接收窗口收缩,压力沿网络反压回客户端。这是特性不是缺陷——broker 过载时保护自己优先存活。
三、RequestChannel:请求队列与请求对象
RequestChannel.Request 携带了处理所需的全部上下文:
- 请求头(api key/version、requestId、clientId、acknowledge 语义);
- 已反序列化的请求体(ProduceRequest/FetchRequest/…);
- 计时埋点:requestQueueTimeNanos(入队时间)、apiLocalTimeNanos(本地处理)、apiRemoteTimeNanos(等待 purgatory/远程)、responseQueueTimeNanos(响应排队)、messageSendTimeNanos(网络发送)。
这套埋点正是 broker 端指标 kafka.network:type=RequestChannelMetrics,name=RequestQueueTime / LocalTime / RemoteTime / ResponseQueueTime 的数据来源(3.x 中由 RequestMetrics 细化),是本讲排查的核心抓手:
- RequestQueueTime 高:IO 线程池不够 / 处理太慢;
- LocalTime 高:CPU 或磁盘同步操作慢(如 index 重建、日志截断);
- RemoteTime 高(尤其 Produce):等 ISR 中 Follower 确认慢——副本/网络问题(呼应第22讲);
- ResponseQueueTime 高:网络线程饱和,回写不动。
四、KafkaRequestHandler:IO 线程池
- num.io.threads(默认 8)个 handler 线程,每个线程主循环就是 requestChannel.receiveRequest() 阻塞式取请求;
- 取到后调用 KafkaApis.handle(request) 按请求类型路由:
| PRODUCE | handleProduceRequest | 等 ISR 确认(RemoteTime) |
| FETCH | handleFetchRequest | 数据可能未就绪 → purgatory |
| LIST_OFFSETS | handleListOffsetRequest | 索引查找,快 |
| METADATA | handleTopicMetadataRequest | 元数据缓存读取,快 |
| JOIN_GROUP/SYNC_GROUP | handleJoinGroupRequest | 协调者延迟操作,组管理 |
| OFFSET_COMMIT/FETCH | consumer offset 读写 | __consumer_offsets 读写 |
4.1 IO 线程里"不能等"的铁律
IO 线程虽比网络线程"重",但依然不能长时间阻塞等待,否则池子被占满、请求队列堆积。所以 Kafka 把所有"可能要等"的步骤都改造为异步 + 回调 + purgatory:
- PRODUCE 等待 acks:消息写入本地日志后,若 acks=all,把"副本确认回调"挂到 ProducerPurgatory,IO 线程立刻返回处理下一个请求;等 ISR 全部 Follower 的 FETCH 拉到该 offset,延迟操作完成,组装响应回写;
- FETCH 等数据:消费者要 offset=100,Leader 日志才到 90,若配置了 fetch.max.wait.ms,把请求挂到 FetchPurgatory,等数据到达(或超时)再响应。
purgatory 的内部实现(层级时间轮)正是第27讲的主题。这里只需要记住:IO 线程只做"快的部分",慢的等待全部丢给延迟队列。
4.2 数据面特例:副本同步 FETCH 的专用通道
Follower 的 FETCH 请求量大且行为固定,Kafka 给它配了专用请求队列与可选的专用处理器(replica.fetcher.threads、以及 ZK 模式遗留的 num.replica.alters 思路;3.x 中主要体现为请求类型优先级与 FetchSession 会话化)。ZK 模式还有 data-plane 与 control-plane 分离;KRaft 模式下所有 broker 请求走 data-plane,controller 的 Raft 通信独立于这套 SocketServer(Raft 有自己的网络通道,raft.message.max.bytes 等)。
五、一次 PRODUCE 请求的完整时序(把所有角色串起来)
ProducerPurgatory
Follower×N
本地日志
Handler-3(IO线程)
requestQueue
Processor-0(网络线程)
Producer
ProducerPurgatory
Follower×N
本地日志
Handler-3(IO线程)
requestQueue
Processor-0(网络线程)
Producer
#mermaid-svg-PSa5Jx4ilQOkhJKN{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-PSa5Jx4ilQOkhJKN .error-icon{fill:#552222;}#mermaid-svg-PSa5Jx4ilQOkhJKN .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-PSa5Jx4ilQOkhJKN .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-PSa5Jx4ilQOkhJKN .marker{fill:#333333;stroke:#333333;}#mermaid-svg-PSa5Jx4ilQOkhJKN .marker.cross{stroke:#333333;}#mermaid-svg-PSa5Jx4ilQOkhJKN svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-PSa5Jx4ilQOkhJKN p{margin:0;}#mermaid-svg-PSa5Jx4ilQOkhJKN .actor{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-PSa5Jx4ilQOkhJKN text.actor>tspan{fill:black;stroke:none;}#mermaid-svg-PSa5Jx4ilQOkhJKN .actor-line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-PSa5Jx4ilQOkhJKN .innerArc{stroke-width:1.5;stroke-dasharray:none;}#mermaid-svg-PSa5Jx4ilQOkhJKN .messageLine0{stroke-width:1.5;stroke-dasharray:none;stroke:#333;}#mermaid-svg-PSa5Jx4ilQOkhJKN .messageLine1{stroke-width:1.5;stroke-dasharray:2,2;stroke:#333;}#mermaid-svg-PSa5Jx4ilQOkhJKN #arrowhead path{fill:#333;stroke:#333;}#mermaid-svg-PSa5Jx4ilQOkhJKN .sequenceNumber{fill:white;}#mermaid-svg-PSa5Jx4ilQOkhJKN #sequencenumber{fill:#333;}#mermaid-svg-PSa5Jx4ilQOkhJKN #crosshead path{fill:#333;stroke:#333;}#mermaid-svg-PSa5Jx4ilQOkhJKN .messageText{fill:#333;stroke:none;}#mermaid-svg-PSa5Jx4ilQOkhJKN .labelBox{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-PSa5Jx4ilQOkhJKN .labelText,#mermaid-svg-PSa5Jx4ilQOkhJKN .labelText>tspan{fill:black;stroke:none;}#mermaid-svg-PSa5Jx4ilQOkhJKN .loopText,#mermaid-svg-PSa5Jx4ilQOkhJKN .loopText>tspan{fill:black;stroke:none;}#mermaid-svg-PSa5Jx4ilQOkhJKN .loopLine{stroke-width:2px;stroke-dasharray:2,2;stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);}#mermaid-svg-PSa5Jx4ilQOkhJKN .note{stroke:#aaaa33;fill:#fff5ad;}#mermaid-svg-PSa5Jx4ilQOkhJKN .noteText,#mermaid-svg-PSa5Jx4ilQOkhJKN .noteText>tspan{fill:black;stroke:none;}#mermaid-svg-PSa5Jx4ilQOkhJKN .activation0{fill:#f4f4f4;stroke:#666;}#mermaid-svg-PSa5Jx4ilQOkhJKN .activation1{fill:#f4f4f4;stroke:#666;}#mermaid-svg-PSa5Jx4ilQOkhJKN .activation2{fill:#f4f4f4;stroke:#666;}#mermaid-svg-PSa5Jx4ilQOkhJKN .actorPopupMenu{position:absolute;}#mermaid-svg-PSa5Jx4ilQOkhJKN .actorPopupMenuPanel{position:absolute;fill:#ECECFF;box-shadow:0px 8px 16px 0px rgba(0,0,0,0.2);filter:drop-shadow(3px 5px 2px rgb(0 0 0 / 0.4));}#mermaid-svg-PSa5Jx4ilQOkhJKN .actor-man line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;}#mermaid-svg-PSa5Jx4ilQOkhJKN .actor-man circle,#mermaid-svg-PSa5Jx4ilQOkhJKN line{stroke:hsl(259.6261682243, 59.7765363128%, 87.9019607843%);fill:#ECECFF;stroke-width:2px;}#mermaid-svg-PSa5Jx4ilQOkhJKN :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
IO线程立即去处理下一个请求
TCP 写入 ProduceRequest
put(Request) 记录入队时间
回到事件循环 继续服务其它连接
take(Request) 记录出队时间
追加消息(PageCache)
acks=all 且ISR未全部确认
挂起 DelayedProduce
FETCH 拉到目标 offset
尝试完成 DelayedProduce
响应放入 Processor-0 的 response 队列
写回 ProduceResponse
实战案例:压测 + 指标观测"排队 vs 处理"
命令演示
# 1. broker 配置建议(server.properties,标注注释)
# num.network.threads=4 # 网络线程:一般 3~8,万兆网卡/多 listener 时上调
# num.io.threads=16 # IO 线程:磁盘数*2 或核数左右起步
# queued.max.requests=1000 # 请求队列:高吞吐场景适当放大,配合监控
# num.replica.fetchers=4 # 副本拉取线程:影响 ISR 跟随速度
# 2. 压测
bin/kafka-producer-perf-test.sh –topic order-events \\
–num-records 5000000 –record-size 1024 –throughput 100000 \\
–producer-props bootstrap.servers=localhost:9092 acks=1
# 3. 观察请求各阶段耗时(JMX,毫秒均值)
# kafka.network:type=RequestChannelMetrics,name=RequestQueueTime, request=Produce
# kafka.network:type=RequestChannelMetrics,name=LocalTime, request=Produce
# kafka.network:type=RequestChannelMetrics,name=RemoteTime, request=Produce
# kafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent
# —— handler 空闲率,持续 <30% 说明 IO 线程吃紧
# 4. 线程栈直查(哪个线程在干什么)
jstack <broker-pid> | grep -A 5 "kafka-request-handler"
jstack <broker-pid> | grep -A 5 "Processor"
小型 Java 演示:从客户端侧测量各阶段延迟分布
import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.Metric;
import org.apache.kafka.common.MetricName;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Map;
import java.util.Properties;
import java.util.TreeMap;
/**
* 第24讲实战:客户端侧观测"请求耗时"与 broker 排队的关联。
* 场景:风控埋点服务发送 20 万条事件,结束后打印生产者自身的
* request-latency 与 record-error-rate 指标,用于和 broker 端指标对照。
*/
public class BrokerLoadDemo {
private static final String BOOTSTRAP = "localhost:9092";
private static final String TOPIC = "risk-events";
public static void main(String[] args) throws Exception {
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.ACKS_CONFIG, "1");
props.put(ProducerConfig.LINGER_MS_CONFIG, 10);
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "lz4");
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 64 * 1024);
int total = 200_000;
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
long start = System.currentTimeMillis();
for (int i = 0; i < total; i++) {
String key = "user-" + (i % 1000);
String value = "{\\"ruleId\\":\\"R-100\\",\\"score\\":0.87,\\"seq\\":" + i + "}";
producer.send(new ProducerRecord<>(TOPIC, key, value), (m, e) -> {
if (e != null) e.printStackTrace();
});
}
producer.flush();
double sec = (System.currentTimeMillis() – start) / 1000.0;
System.out.printf("发送 %d 条耗时 %.1fs,吞吐 %.0f 条/s%n", total, sec, total / sec);
// 打印与"请求-响应链路"相关的客户端指标,与 broker 端指标对照
Map<String, Metric> sorted = new TreeMap<>();
for (Map.Entry<MetricName, Metric> e : producer.metrics().entrySet()) {
String name = e.getKey().name();
if (name.contains("request-latency") || name.contains("requests-in-flight")
|| name.contains("record-error-rate") || name.contains("buffer-available-bytes")) {
sorted.put(e.getKey().group() + "." + name, e.getValue());
}
}
sorted.forEach((k, v) ->
System.out.printf("%-55s = %s%n", k, v.metricValue()));
}
}
}
关键点解读
踩坑提示 / 生产建议
本讲小结
broker 的请求处理是标准的 Reactor 分层模型:Acceptor 接连接、Processor 网络线程收发字节并维护响应队列、共享请求队列削峰、KafkaRequestHandler 池执行业务逻辑,慢等待全部交给 purgatory。四段式耗时指标(排队/本地/远端/响应排队)是定位性能问题的手术刀。线程参数调整务必以指标为依据,而不是拍脑袋加倍。
思考题