Apache Beam 2.2.0 Dataflow管道报错:拒绝拆分,作业失败求助
首先还原你遇到的错误日志:
拒绝在ShufflePosition(base64:AAAAAtxW0XsAAQ)处拆分<位于混洗范围[ShufflePosition(base64:AAAAAgD_AP8A_wD_AAE), ShufflePosition(base64:AAAAAtxW0XsAAQ))的ShufflePosition(base64:AAAAAtxW0XoAAQ)位置>:提议的拆分位置超出范围。作业最终失败,作业ID:2018-04-02_00_19_15-14115706867296503746。
这个问题我之前处理过,本质是Apache Beam 2.2.0版本中reshuffle操作的一个已知bug——早期版本的reshuffle在计算Shuffle分片的边界范围时存在逻辑错误,当Dataflow尝试对Shuffle数据进行拆分调度时,会错误地计算出超出合法范围的拆分位置,最终导致作业崩溃。
给你两个可行的解决思路:
优先升级Beam版本
这个bug在Beam 2.3.0及之后的版本已经被完全修复了。如果你用Maven管理依赖,直接修改pom.xml中的Beam核心和Dataflow运行器版本即可:<!-- Beam核心依赖 --> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-sdks-java-core</artifactId> <version>2.3.0</version> <!-- 推荐用更高的稳定版本,比如2.x系列的最新版 --> </dependency> <!-- Dataflow运行器依赖 --> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-runners-google-cloud-dataflow-java</artifactId> <version>2.3.0</version> </dependency>升级版本不仅能解决这个Shuffle拆分问题,还能获得后续版本的性能优化和其他bug修复,是最彻底的解决方案。
临时规避方案(无法立即升级时)
如果暂时没法升级版本,可以尝试调整reshuffle的使用方式:- 避免在超大数据集后直接使用
reshuffle,可以先通过GroupByKey做一次轻量聚合,再执行reshuffle操作,减少单个Shuffle分片的数据量; - 手动指定reshuffle的分片数,通过
withNumShards()设置合适的并行度,比如:// 根据你的数据规模调整分片数,比如64、128等 pipeline.apply(Reshuffle.viaRandomKey().withNumShards(64));
这种方式能降低触发边界计算错误的概率,但只是临时手段,还是建议尽早升级版本。
- 避免在超大数据集后直接使用
内容的提问来源于stack exchange,提问作者Ashley Thomas

