GCP Dataflow批处理作业Reshuffle步骤解码失败问题求助
GCP Dataflow批处理作业Reshuffle步骤解码失败问题求助
大家好,我最近在调试GCP Dataflow批处理作业时碰到了一个棘手的问题——作业在Reshuffle()步骤直接失败,抛出了解码相关的ValueError,具体错误信息如下:
ValueError: Error decoding input stream with coder WindowedValueCoder[TupleCoder[LengthPrefixCoder[DeterministicFastPrimitivesCoder], StateBackedIterableCoder[LengthPrefixCoder[FastPrimitivesCoder]]]]
跟大家说下我的作业逻辑:我们有一个步骤会生成若干个列表对象,之后需要把这些列表传递给下一个PCollection做后续处理。最开始没加Reshuffle()的时候,这两个步骤被Dataflow自动融合了,导致列表只能线性处理,没法利用多worker并行执行。为了解除融合、把列表分配到不同worker上并行处理,我才添加了Reshuffle()步骤,但现在在Reshuffle()内部的GroupByKey环节触发了上面的解码错误。
我已经尝试过一个方案:把第一个PCollection的输出从单纯的列表对象改成包含列表的自定义对象,但还是遇到了完全相同的解码错误,实在找不到头绪了,想请教下社区的各位大佬,有没有遇到过类似的问题?或者有什么排查方向可以给我参考?
备注:内容来源于stack exchange,提问作者Ben Delany
相关产品推荐
相关产品推荐

