DataFlow中Beam流作业读取Pub/Sub疑似卡顿问题咨询
可能成因
- Pub/Sub读取端Watermark追踪机制瞬时过载
Beam的Pub/Sub IO读取层需要持续追踪所有订阅分区下未确认消息的最早事件时间,以此计算全局Watermark。当吞吐量从25k/s瞬间跳涨到130k/s时,待追踪的未ack消息元数据量会在短时间内膨胀数倍,Watermark聚合计算线程会出现短暂阻塞,无法正常推进Watermark;等元数据批量同步完成时,Watermark会出现跳变式拉升也就是观测到的尖刺。这个阻塞过程中读取线程无法正常拉取消息,就会形成短暂的消息积压,等计算恢复后积压会快速消化。 - 固定Worker池的线程调度抖动
固定规模集群虽然总资源余量足够承载高峰负载,但每个Worker的执行线程池容量是固定值。高峰流量下下游窗口聚合、分组操作会瞬间占用大量工作线程,Pub/Sub读取侧的消息拉取、ack确认、Watermark上报这几个核心线程会被操作系统调度暂时抢占优先级,出现数秒的执行停顿。停顿累积2-3分钟就会触发Dataflow的Possible stuckness阈值告警——这个告警只基于阶段进度上报间隔触发,不代表作业真的卡死,调度恢复后进度追平告警就会自动消失,也不会在业务日志里留异常记录。 - 下游Shuffle反压的传导延迟
负载升高叠加分组聚合触发内部Shuffle的推测是成立的:聚合操作触发Shuffle写入时,Shuffle缓冲区会在短时间被打满,反压信号从聚合节点传递到Pub/Sub读取节点存在10~30秒的延迟。这段时间读取端还在持续拉取消息,等反压信号传导到读取端时,拉取动作会被临时暂停,没有新的事件时间输入时Watermark会跟随系统时间向前跳变形成尖刺;等Shuffle缓冲区把积压的批次数据写完,反压解除,读取端恢复拉取,Watermark就会回落到真实的消息时间位置,积压也会快速消化。
排查方向
- 对齐时间轴拉取Dataflow阶段级监控指标,重点核对三个指标和Watermark尖刺的时间重合度:如果尖刺时段Pub/Sub ack请求延迟同步升高,但拉取请求数没有明显下降,基本可以判定是Watermark元数据追踪过载;如果尖刺时段拉取请求数直接跌到0,就是被下游反压或者线程调度抢占导致的。
- 检查Worker节点的系统层指标:如果尖刺时段用户态CPU占比没有打满,但系统态CPU占比、CPU运行队列长度突然升高,就是线程调度抖动导致的。可以适当调大每个Worker的执行线程数冗余,或者把Pub/Sub读取的Watermark更新间隔从默认100ms调整到500ms,降低元数据计算的开销。
- 查看下游聚合阶段的Shuffle指标:如果尖刺出现前10~30秒,Shuffle缓冲区占用刚好到达配置阈值,就是反压传导导致的尖刺。可以适当调大Shuffle内存缓冲区的占比,或者在分组聚合前加本地预聚合逻辑,减少Shuffle阶段的瞬时传输数据量,削平流量峰值。
- 检查Pub/Sub订阅的分区负载分布:如果高峰时段存在少数热点分区(单分区流量占总流量30%以上),会导致对应读取Task的局部负载过高,拉取、ack、Watermark计算都卡在少数Worker上,也会触发短暂的Watermark尖刺。这种情况可以适当增加Pub/Sub主题的分区数,开启消息键的分区均衡策略即可缓解。
注:如果Watermark尖刺可以快速回落、消息积压可以持续消化、没有出现端到端延迟无限上涨的情况,这类瞬时抖动属于高负载下的正常现象,
Possible stuckness属于阈值触发的提示类告警,不代表作业存在逻辑故障。
内容的提问来源于stack exchange,提问作者Sergii V.
相关产品推荐
相关产品推荐

