ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka与RocketMQ长轮询机制深度解析:Purgatory、HoldService与生产调优

Kafka与RocketMQ长轮询机制深度解析:Purgatory、HoldService与生产调优 前几天群里有人问Kafka 的fetch.max.wait.ms调成多少合适我反问了一句这个参数是谁在等待等待的时候服务端到底在干什么对面半天没吱声。说实话RocketMQ 和 Kafka 虽然都在讲“长轮询”但很多人把二者混为一谈以为都是“拉不到数据就挂起一段时间”。可挂起谁、等多久、由谁唤醒、唤醒后重新走什么路径这两家的实现差异非常大。这篇文章想把这层窗户纸捅破。我会结合源码层面的行为分五块讲清楚先看为什么 RocketMQ 和 Kafka 最终都选择了 Pull 长轮询的组合然后分别拆 Kafka 和 RocketMQ 的等待机制接着把两者放到同一张对比表里梳理参数、唤醒时机、超时模型和真实场景最后聊生产环境里我实际踩过的坑和调优手段以及如何确认长轮询确实在按预期工作。1. 从短轮询到长轮询为什么消费端非要等1.1 名义上的 Push 和实际上的 Pull先消除一个最常见的技术误解RocketMQ 的 PushConsumer 和 Kafka 的KafkaConsumer.poll()一个看起来像“推”一个看起来像“拉”但本质上都是消费者主动拉取。RocketMQ 所谓的 Push只是客户端把拉取循环、消息回调、offset 管理全部封装好了你只需要注册一个MessageListener。底层依然是PullMessageService线程不断构造PullRequest去 broker 拉消息拉回来之后回调消费消费完成再发下一轮。为什么两大消息中间件最终都选了 Pull 而不是服务端主动 Push有三个原因特别朴素顺序性Push 模型下broker 不知道消费者当前处理到哪一条一旦堆积链路缓冲区里全是消息顺序难以保证。Pull 模型里消费者拉多少处理多少天然维护顺序。背压控制消费者处理慢就少拉一点处理快就多拉一点broker 不需要感知每个消费者的 CPU、GC 和线程池状态。状态管理broker 不需要为每个消费者维护推送会话和游标状态消费进度的维护集中到了客户端和 offset 存储服务端实现大幅简化。顺着这个思路往下走拉取模型就暴露出一个明显的短板怎么知道“现在有数据了”如果调用方干等延迟高如果疯狂轮询浪费资源。长轮询正是为了解决这个矛盾被发明的。1.2 短轮询的痛点空请求与无效开销在长轮询普及之前很多客户端实现的是短轮询消费者发送 fetch 请求broker 检查发现没有新数据立刻返回一个“空”。消费者收到空结果休息一小段间隔比如 50ms 或者 100ms然后再重试。这个模式的问题非常直接空轮询率极高。topic 不活跃时每个消费者每 100ms 发一个请求服务端每次都要做一次全链路检查然后返回空。10 个消费者、20 个分区每秒的空响应就有上百个。延迟不稳定。新消息如果刚好落在两次轮询之间消费者最长要等一个完整的轮询间隔才能看到消息。这个间隔本身就是消费延迟的下限。请求风暴。broker 重启、分区 leader 切换时消费者会一起重连、一起拉取没有数据又一起开始短轮询服务端瞬时请求量直接冲高。我在测试环境做过实验一个没有任何消息的 topic单消费者短轮询间隔 50msbroker 每秒收到的 Fetch/Pull 请求大约 20 个改成相同超时参数的长轮询后请求量直接下降两个数量级。长轮询表面上只是“少发了几次请求”实际对 broker 的连接数、请求处理线程和 GC 压力都有明显改善。1.3 长轮询的通用模型挂起请求不是挂起连接长轮询的通用做法可以概括为一句话请求到达 broker 后如果暂时没有数据可返回不立即返回空而是把这次请求作为“延迟任务”记下来等条件满足后再补响应。条件通常有三类有新数据可拉取或已选数据的字节数达到阈值。等待超时兜底返回一个空响应。请求被取消比如客户端断开连接。这个过程很像去餐厅吃饭。菜没好时你不会每隔十秒钟跑去后厨催一次而是告诉服务员“好了喊我”。服务员会在菜好的那一刻来通知你。这里要特别强调一个关键点挂起的是请求不是连接。很多人误以为长轮询是把 socket 连接一直占着、服务端线程阻塞在那里。实际上Kafka 和 RocketMQ 的挂起都是在请求处理逻辑内部实现的延迟判定处理线程返回一个“稍后响应”的标志请求对象被保存到内存结构中线程立刻去处理下一个请求。真正把响应发回 socket要等延迟任务被触发。这个认知是理解后续所有源码细节的前提。2. Kafka 的长轮询实现Purgatory 里的 DelayedFetch2.1 客户端侧poll 并不是轮询Kafka 消费端入口是KafkaConsumer.poll()名字听起来像轮询但真正拿消息的过程由后台 fetcher 线程负责。每次poll()触发时如果本地缓冲没有足够可返回的记录客户端就会向分区 leader 所在的 broker 发送 FetchRequest。决定这个请求会不会“空等”的主要是请求头里的两个参数。第一个是fetch.min.bytes默认 1 字节。它表示本次 fetch 响应至少需要多少字节。如果 broker 上可用数据不足这个值broker 不会立刻返回而是进入等待。第二个是fetch.max.wait.ms默认 500。它表示在未达到fetch.min.bytes的情况下broker 最多等待多久。超时后即使字节数不够也把已攒到的数据返回。所以fetch.max.wait.ms不是“轮询间隔”而是“攒批的耐心上限”。想要低延迟就把fetch.min.bytes调小、fetch.max.wait.ms调小想要高吞吐、减少请求数就调大fetch.min.bytes让 broker 攒够一批再返回。这里还要提一下fetch.max.bytes和max.poll.records的分工前者限制单个 FetchRequest 响应体的最大字节数避免一次拉太多导致客户端 OOM后者限制poll()返回给用户的最大记录条数。这两个参数和长轮询的等待条件没有直接关系但共同决定了消费端实际拿到的批大小。2.2 服务端核心DelayedOperationPurgatory 与 DelayedFetchKafka broker 收到 FetchRequest 后会先尝试直接处理从本地日志中算出每个分区可返回的 offset 范围和数据。如果满足fetch.min.bytes立即组装响应返回如果不满足请求就会被包装成一个DelayedFetch放进DelayedOperationPurgatory等待。Purgatory是 Kafka 实现延迟操作的组件底层采用时间轮管理超时。每个 DelayedFetch 注册两类回调超时回调时间轮扫描到过期任务如果实在等不到足够数据就返回当前已拿到的数据或者返回空。数据到达回调分区日志有新消息追加、HW 推进时触发tryComplete()重新计算当前所有分区现在是否满足返回条件。正是因为采用“时间轮 回调”而不是“线程阻塞 等待”Kafka 的 broker 线程在请求等待期间不会被占用。单个 broker 可以同时挂起成千上万个 fetch 等待任务内存也能保持可控。刚开始读这块代码时我有一个误解以为每次有新消息进来挂在请求队列里的 fetch 都会被全部唤醒一遍。实际并不是。tryComplete()会检查当前攒到的总字节数是否已经满足fetch.min.bytes没满足就主动放弃完成让请求继续挂在时间轮里。所以 Kafka 的长轮询不是“来一条推一条”而是“攒到一定量再返回”。这一点决定了它的吞吐模型和 RocketMQ 完全不同。2.3 水位推进与事务消息对长轮询的影响Kafka 长轮询还有一个隐形门槛服务端计算“可返回数据”时不是直接看日志的 LEO而是看消费者可见的 HW。也就是说如果分区 follower 副本同步落后leader 的 HW 不前进消费者即使使用长轮询也只能看到 HW 之前的位置。很多同学排查消费延迟时只看 consumer lag不看 HW。结果 partition 的 LEO 明明很高消息却一直消费不到长轮询反复返回空最后定位下来是某个 follower 副本卡住导致 HW 不推进。事务消息同理。未提交的事务消息默认不会返回给消费者只有事务 commit 之后LSO 推进消息才变得可见。如果你用的是read_committed隔离级别长轮询对“可拉取 offset 范围”的判定会更严格。这两点和长轮询的关系很紧密因为新消息写入并不代表立即触发唤醒只有消息“对消费者可见”才会。所以排查 Kafka 消费延迟时除了看请求参数还要看 HW 是否在正常推进。3. RocketMQ 的长轮询实现PullRequestHoldService 的 15 秒等待3.1 客户端拉取循环PullMessageService 与 ProcessQueueRocketMQ 的 PushConsumer 启动以后后台会有一个PullMessageService线程持续构建PullRequest投递到内部的BlockingQueuePullRequest。很多资料在这里就直接跳到“长轮询”了其实有个细节很关键PullRequest不是发完就结束。每次从 broker 拉回消息后回调逻辑会检查消费状态再生成一个新的PullRequest塞回队列等待下一轮。正是因为每一轮拉取完成后自动续上下一轮客户端看起来才像“推模式”。整个过程中ProcessQueue负责维护消息队列的滑动窗口。broker 拉回的消息先放进ProcessQueue由消费线程池取走消费消费完成后再更新消费位点生成下一轮PullRequest。pullBatchSize默认是 32 条也就是单次拉取最多 32 条。这个批大小比 Kafkamax.poll.records的默认值500小很多所以 RocketMQ 在消息稀疏时延迟更低但同样的数据量下客户端处理循环会更频繁。3.2 Broker 侧挂起PullRequestHoldService现在看 broker 侧。客户端发来的拉取请求由PullMessageProcessor处理。处理过程大致是先从ConsumeQueue查当前消费位点之后有没有消息有就直接返回如果没有数据同时配置项longPollingEnabletrue默认是 true就把请求交给PullRequestHoldService继续等待。PullRequestHoldService内部有一个以topicconsumerGroup为 key 的 Mapvalue 是ManyPullRequest。ManyPullRequest内部维护了一个ArrayBlockingQueuePullRequest保存所有挂起请求。后台线程循环执行两类动作定时扫描所有挂起的请求收到messageArriving信号时立刻去遍历对应队列的挂起请求。默认的挂起超时是longPollingTimeout在 BrokerConfig 里默认值是 15000ms。也就是说一个请求挂满 15 秒还没有等到新消息HoldService 就会创建空响应返回客户端收到空结果后重新发起下一轮。这里和 Kafka 的差异就很明显了Kafka 默认 500ms 就兜底返回RocketMQ 默认 15 秒。但 RocketMQ 并不依赖这个 15 秒来保证实时性因为新消息到达时立刻会唤醒挂起请求15 秒只是极端场景下清理空挂请求用的。3.3 新消息到达时如何唤醒挂起请求重点看这条唤醒链路。RocketMQ 的DefaultMessageStore在消息落盘后会通过doDispatch把消息分发到 ConsumeQueue分发过程中会回调messageArriving。PullRequestHoldService收到通知后按照消息的topic和queueId找到对应的挂起请求列表逐个唤醒。这里有一个很微妙的细节唤醒条件是“该消费队列位点之后有新消息可消费”而不是“物理日志写入就算”。检查时会读取ConsumerOffsetManager里的消费进度再和 ConsumeQueue 当前最大 offset 比较。如果消费者没有及时上报消费进度即使物理日志写入了新消息挂起的请求也可能不满足唤醒条件继续等到超时。所以 offset 上报节奏对长轮询的命中率有直接影响。唤醒之后请求重新回到PullMessageProcessor处理流程再次走常规查询。如果有消息就直接组装响应返回如果仍查不到数据就再次挂回PullRequestHoldService。客户端不需要感知这次“假唤醒”看起来只是响应返回得稍晚了一些。3.4 RocketMQ 长轮询的常见误解一个高频误解RocketMQ 把消费者的连接挂起 15 秒。实际上服务端挂起的是PullRequest对象不是 Netty Channel。Netty 线程处理完这个请求后马上会去处理其他请求真正等待的是PullRequestHoldService的后台线程。另一个误解挂起 15 秒内一定会等到消息所以消费延迟不会超过 15 秒。准确说如果新消息一直不来消费者要等满 15 秒才会收到空响应如果新消息在超时前 1 秒到达唤醒后响应很快回来实际延迟不到 1 秒。只有在唤醒信号和后台扫描发生临界竞争时延迟才可能接近 15 秒。这个概念对后续调优很重要。4. 双雄对比从实现差异看设计哲学4.1 拉一张直给的对照表对比维度KafkaRocketMQ核心等待参数fetch.min.bytes/fetch.max.wait.mslongPollingTimeout默认 15s数据到达唤醒分区 HW 推进后触发DelayedFetch.tryComplete()消息 dispatch 到 ConsumeQueue 后触发messageArriving空数据兜底时间500msfetch.max.wait.ms默认值15000mslongPollingTimeout默认值请求挂起实现DelayedOperationPurgatory时间轮管理PullRequestHoldService后台线程 阻塞队列挂起粒度一个请求可覆盖多个分区按总字节数判断一个请求对应一个 topicqueueId按 Offset 判断返回条件未达到 min.bytes 就继续等超时才返回新消息到达立即返回超时兜底批量控制侧重面向字节数min/max bytes面向条数pullBatchSize默认 32典型延迟模型攒批模型延迟换吞吐快速响应模型优先低延迟这张表可以直接拿去当面试总结用也可以作为线上调优的对照基准。两张表里唯一需要反复强调的点是Kafka 等待的是“累计字节数”RocketMQ 等待的是“当前消费队列有没有新 offset”。4.2 为什么 Kafka 攒批RocketMQ 偏向快速返回Kafka 的一个 FetchRequest 可以同时覆盖多个分区broker 等待时累计的是所有分区当前位置之后的字节总数达到fetch.min.bytes才返回。这个设计天然假定“这一批数据值得等待”因为一次网络传输、一次响应解析的成本是固定的攒成大批次可以摊薄这些固定开销。所以 Kafka 用户调高吞吐时会刻意把fetch.min.bytes调到几十 KB 甚至 1MB让 broker 多攒一会儿。RocketMQ 的 PullRequest 只针对单个消息队列不存在“多个队列凑字节数”的需求。它的目标很纯粹有消息马上回没消息挂一会儿。所以默认挂起超时可以给到 15 秒而不担心延迟因为新消息到达时有独立唤醒路径优先级更高。这也解释了网上经常吵的话题“RocketMQ 延迟比 Kafka 低”。这个说法不绝对。单分区单消息、空闲等待场景下RocketMQ 确实能通过“来一条唤醒一次”实现更低延迟但在大量并发消息场景下Kafka 的攒批能用大包摊销网络成本吞吐更稳。讨论延迟之前先确认消息模型和压力模型否则结论没有意义。4.3 对连接数和请求槽位的影响长轮询不占用专用线程但会占用“未完成请求”的槽位。Kafka 客户端与 broker 之间有多个连接同一连接上可以同时存在多个未完成的 fetch 请求吗正常情况下可以但消费端通过max.in.flight.requests.per.connection控制在途请求数。如果fetch.max.wait.ms设置过大消费者在等待响应期间这个连接上的后续新请求就会被排队间接影响其他 topic 的拉取。RocketMQ 的请求模型是“同一连接、同一时间只处理一个请求”PullRequest 串行执行。所以消费者线程数不等于并发挂起数。每个 broker 上同一消费组对同一 queue 只可能有一个挂起的 PullRequest。这个模型更简单但也意味着消费端处理慢时拉取请求会被拖住形成“拉取耗时大但消费没有堆积”的假象。我在生产环境见过很多次类似问题最终都是把拉取循环和消费线程池分离后解决的。5. 生产环境长轮询调优参数、监控、坑与心得5.1 参数配置速查可以直接抄作业Kafka 低延迟场景交易通知、订单状态同步fetch.max.wait.ms100~300fetch.min.bytes1max.poll.records100~200使用手动提交enable.auto.commitfalseKafka 高吞吐批处理场景日志传输、数仓同步fetch.min.bytes64KB~1MBfetch.max.wait.ms500~5000fetch.max.bytes50MB可适当调大num.consumer.fetchers增加后台拉取线程RocketMQ 常规场景挂起超时longPollingTimeout保持默认 15s不用刻意调小真正影响延迟的是唤醒链路pullBatchSize默认 32大消息建议调小小消息可以调到 64~128消费线程池consumeThreadMin/consumeThreadMax与拉取频率配合避免消费侧成为瓶颈如果对 RocketMQ 延迟要求更苛刻可以调小客户端的pullInterval默认 0 表示拉完立即接着拉但这会增加 broker 压力空置 topic 没必要这样设置。5.2 我真实踩过的几个坑第一个坑Kafka 的max.poll.interval.ms和长轮询的组合问题。如果消费者处理一批消息的时间接近max.poll.interval.ms默认 5 分钟及时 fetch 线程正常消费组也可能判定消费者失联而触发 rebalance。我遇到过一次离线任务consumer 有长时间 GCfetch 线程正常但 poll 主线程没跟上消费组反复重平衡消费进度一直倒退。排查时一定要分清楚是拉取侧慢还是消费侧慢不能只看 fetch latency。第二个坑RocketMQ 唤醒路径对消费进度的依赖。前面提到messageArriving要对比 ConsumeQueue 最新 offset 和消费者上报 offset。如果消费者上报延迟很大新消息写入后不会立刻唤醒挂起的 PullRequest表现为“明明有消息延迟却很高”。生产环境建议关注 broker 端 offset 上报节奏合理设置autoCommitInterval。第三个坑fetch.max.wait.ms不是越大越好。一位同事为了减少 broker 请求量把 Kafka 的fetch.max.wait.ms调成 10000结果消费延迟平均值直接升到 10 秒级别。原因是 topic 消息量太小永远达不到fetch.min.bytes每个分区实打实等满 10 秒。低流量 topic 的fetch.max.wait.ms保持在 500 以内更合理不要用拉长等待时间的方式换请求量。第四个坑RocketMQ 长轮询挂起数量与 broker 线程的关系。当同一 broker 挂起的 PullRequest 很多时PullRequestHoldService的后台线程遍历 Map 的开销会变大。更麻烦的是唤醒时的“惊群”效应一个 topic 来消息会唤醒该 topic 下所有消费组的挂起请求每个都要重新走一遍查询。高频 topic 且有大量消费组订阅时关注 broker 日志里PullRequestHoldService的执行耗时必要时降低longPollingTimeout用更频繁的空返回打断长挂起。第五个坑Kafka 请求超时和长轮询的边界。消费者发送 fetch 后broker 最多等到fetch.max.wait.ms但 broker 端还有一个统一的request.timeout.ms默认 30 秒限制整个请求周期。如果fetch.max.wait.ms设置值超过request.timeout.ms会直接报超时。这个网上资料很少提一旦你往大了调fetch.max.wait.ms就很容易碰到。5.3 怎么确认长轮询真的在生效Kafka 有两个直接指标。消费端看kafka.consumer:typeconsumer-fetch-manager-metrics里的fetch-latency-avg如果这个值接近fetch.max.wait.ms说明大部分 fetch 是空转等待如果远小于该值说明数据充足请求很容易被直接满足。broker 端看kafka.network:typeRequestMetrics,nameRequestsPerSec,requestFetch空闲时每秒请求数很低说明长轮询把空轮询抑制住了这就是正常状态。RocketMQ 侧可以看 broker 日志里PullRequestHoldService的请求释放日志也可以用 JMX 观察DefaultMessageStore的dispatchBehindBytes如果该值接近 0说明唤醒链路没有积压。抓包是最直接但最有效的验证方式订阅一个空 topic抓 broker 端口的请求记录长轮询下一次请求到响应之间的间隔可以达到接近超时值短轮询则是固定的短间隔高频往返。5.4 最后分享一点个人经验调长轮询从来不是单独调一个参数就能完成的事情。Kafka 要把fetch.min.bytes、fetch.max.wait.ms、max.poll.records连在一起看RocketMQ 要把pullBatchSize、pullInterval、消费线程池和 offset 上报节奏连在一起看。参数之间互相耦合只压一个指标很容易顾此失彼。我自己的习惯是先画一条“消息产生到消费完成”的链路分三段看broker 写入段、拉取等待段、消费处理段。哪一段耗时最长就先查哪一段。长轮询属于拉取等待段它只解决“拉取不空转”的问题如果写入段有堆积或者消费段处理太慢单纯调长轮询参数根本看不出效果。另外有个小技巧压测消息队列时不要只测满负载吞吐一定要测“稀疏消息 空转”场景下的延迟和请求量。很多线上事故都出在低负载时的异常请求模式上。长轮询在满负载下的表现反而不容易出问题真正体现调优功底的是它能不能在空闲时安静地等待、在消息到达的一瞬间又快又准地醒来。
返回列表