Dataflow作业CPU使用率极低却无法追上Pub/Sub消息积压该如何排查?
配置问题排查
- 版本兼容缺陷
你当前使用的apache-beam==2.32.0是2021年发布的老旧版本,Python流式处理的Pub/Sub读取优化、Dataflow Prime兼容性适配都存在大量已知问题,优先升级到2.40.0及以上的稳定版本,同步确认Python版本匹配即可。 - 小流量场景参数缺失
Dataflow默认参数面向高吞吐量场景优化,会默认攒够一定数量的消息才触发处理,对于你每秒1条的极低流量场景,需要补充以下启动参数:
--experiments=min_bundle_size=0 --experiments=max_bundle_time_millis=100
同时给ReadFromPubSub步骤增加拉取参数,降低拉取间隔:
beam.io.gcp.pubsub.ReadFromPubSub( subscription=known_args.input_subscription, timestamp_attribute=None, with_attributes=True, max_messages_per_pull=1000, pull_retry_delay_seconds=0.1 )
- 扩缩容策略适配问题
你当前使用THROUGHPUT_BASED自动扩缩容,在极低吞吐量场景下不会触发扩容决策,同时初始2个worker会产生不必要的跨worker协调开销,可先调整为单worker测试,观察处理速度变化。
日志排查方向
- 优先核对核心作业指标
进入Dataflow作业详情的「指标」页,优先查看三类指标:- Pub/Sub订阅的未确认消息堆积量,如果持续上涨说明读取环节速度跟不上
ReadFromPubSub步骤的每秒处理消息数,如果长期低于1条/s,可确认问题出在消息拉取环节- 各步骤的平均处理延迟,定位延迟最高的环节
- 筛选worker日志排查异常
在日志筛选器中选择「worker日志」,过滤日志级别为WARNING、ERROR,搜索以下关键词:Pubsub:检查是否存在拉取超时、权限报错、重试次数超限等问题Throttling:检查是否存在Pub/Sub限流、Dataflow资源限流的提示Bundle:检查是否存在数据处理包处理失败重试的记录,重试间隔过长会大幅拉低整体处理速度
- 增加自定义日志定位耗时
可在CustomParsing的process方法中增加自定义INFO日志,打印消息的Pub/Sub发布时间、处理时间的差值,即可明确耗时是出现在Pub/Sub拉取环节还是后续处理环节。
内容的提问来源于stack exchange,提问作者kiki
相关产品推荐
相关产品推荐

