Google Dataflow Worker卡在99%完成状态的技术求助
Dataflow管道99%进度无限挂起问题求助
困扰已超过一周:运行Dataflow数据管道时,Worker会处理完绝大多数数据,但到99%进度时便无限挂起。应用日志显示最后一个输入已被拖慢的Worker成功转换(并非卡在Map阶段),无法确定具体原因。
希望为整个任务添加超时机制,终止挂起的Worker并利用现有数据完成管道,但Python的超时库局限性极强(无法使用signal,func_timeout仅支持非嵌套场景,无法创建子进程等)。
相关情况
- 拖慢的Worker:仅单个Worker出现拖慢情况,其余Worker已完成任务
- 管道阶段:任务包含ReadFromText、两次Reshuffle、FlatMap转换、WriteToTFRecord等阶段
- 受影响阶段的影响:当吞吐量下降时任务本应结束(实际已完成99.9%),但仍需为Worker空转的时间付费
- 补充说明:"train"阶段处理的数据量比Val/Test阶段多一个数量级
已尝试的修复方案
- 为处理函数添加超时:提升了吞吐量,但问题仍存在
原因判断
认为是Dataflow端的偶发Bug导致管道异常(可能是资源不足、工作窃取时死锁、网络中断等),因此更关注如何提升管道的健壮性,而非预判并阻止所有可能的异常。
管道代码(Train/Val/Test三个阶段代码结构一致)
"ReadFromText" >> beam.io.ReadFromText(file_pattern=data_path) | "Shuffle0" >> beam.Reshuffle() | "Convert" >> beam.FlatMap(map_fn) | "Shuffle1" >> beam.Reshuffle() | "Write" >> beam.io.WriteToTFRecord( output_path, coder=beam.coders.ProtoCoder(tf.train.Example)))
内容的提问来源于stack exchange,提问作者Rajan Saini
相关产品推荐
相关产品推荐

