分支PCollection融合阶段的重试处理机制是怎样的?
带分支融合阶段的Dataflow Bundle重试机制
Apache Beam运行器默认会对执行图做fusion(融合优化),如果不主动插入断开融合的操作,多个关联的PTransforms会被合并为同一个执行Stage(阶段)。在Cloud Dataflow中,原本用来断开融合的Reshuffle转换已被标记废弃,运行时内置的ReshuffleOverrideFactory会把它自动简化为「带窗口的GroupByKey+可迭代对象展开」的组合操作,和直接用GroupByKey断融合的效果一致。
按照官方故障处理规则,Cloud Dataflow运行过程中如果出现故障,会重试失败的bundle(数据束):流处理模式下无限重试,批处理模式下最多重试4次。Beam模型默认逻辑是仅重试失败的转换,只有耦合失败(Coupled Failing)场景下会触发关联步骤重试。
对照测试场景
测试构造了两个结构相近的管道,统一配置为向不存在的Pub/Sub主题写入以触发错误:
- 管道A:
Read Files >> Parse Files (High Fan-out ParDo transform) >> [Write to BigQuery, Write to Pub/Sub] - 管道B:
Read Files >> Parse Files (High Fan-out ParDo transform) >> [Reshuffle >> Write to BigQuery, Write to Pub/Sub]
管道A观测结果
- 负责文件解析的
ParDo步骤因为Pub/Sub写入失败被持续重试 - 系统延迟持续升高,全程无数据写入BigQuery:因为BigQuery写入连接器内部自带
Reshuffle,存在错误元素的bundle会触发同一个融合阶段内,Reshuffle的groupbykey操作之前所有步骤全部重试,和预期逻辑一致。
管道B观测结果
- 负责文件解析的
ParDo步骤没有出现重试 - 虽然Pub/Sub指标上报了写入错误(按常规认知应该重试整个文件处理流程、不会有数据写入BigQuery),但实际观测到数据成功写入了BigQuery。
核心规则说明
两个管道表现差异的核心原因是:融合边界就是重试隔离的边界,重试范围永远不会跨融合阶段。
- 同一个融合阶段内的所有步骤(包括多分支输出的所有分支)是强耦合执行的:只要阶段内任意一个分支的处理抛出错误导致bundle失败,整个阶段的bundle执行结果都会被丢弃,下次重试会从头执行这个阶段内所有分支的逻辑。这就是管道A的表现来源:Parse ParDo、Write to Pub/Sub、BigQuery写入自带Reshuffle之前的逻辑全在同一个融合阶段,Pub/Sub写入失败会导致整个阶段的所有输出都不生效,Parse步骤被反复重试,BigQuery也拿不到有效输入。
- 一旦插入
Reshuffle(也就是运行时被替换成的GroupByKey类shuffle操作),就相当于在这个位置切出了新的融合边界,边界两侧属于完全独立的执行阶段,重试逻辑互不干扰。管道B里Reshuffle把BigQuery写入分支和上游Parse、Pub/Sub写入分支拆成了两个独立阶段:Parse和Pub/Sub写入同属一个上游阶段,这个阶段里Pub/Sub写入失败确实会触发该阶段的bundle重试,但重试过程中只要Parse步骤成功输出了元素,这些元素在经过Reshuffle对应的shuffle流程持久化之后,就属于下游独立阶段的输入了。上游阶段的Pub/Sub分支后续再报错重试,不会影响已经被shuffle持久化、交给下游BigQuery写入阶段的数据,所以才会出现Pub/Sub报写入错误,但BigQuery仍然能成功收到数据的现象;Parse步骤没有被重试,也是因为它输出到Reshuffle的元素已经被持久化确认,不需要重复计算。
总结:发生bundle重试时,重试范围是单个融合阶段内的全部逻辑。如果一个融合阶段存在多个输出分支,只要阶段内任意分支报错导致bundle失败,该阶段所有分支的当次执行结果都会被丢弃,下次重试会重新执行整个阶段的全部分支逻辑;被shuffle类操作形成的融合边界隔开的其他阶段,不会被连带重试。
内容的提问来源于stack exchange,提问作者Seng Cheong
相关产品推荐
相关产品推荐

