求助解决CoGroupByKey中keyCoder非确定性异常问题
解决Apache Beam CoGroupByKey的非确定性KeyCoder异常
问题背景
需要在两个PCollection之间通过id1和id2键值对执行CoGroupByKey操作,已定义包含这两个键的自定义类KeysId1Id2,并实现了对应的Coder类KeysId1Id2Coder按id1→id2的固定顺序编码,但始终抛出以下异常:
java.lang.IllegalStateException: the keyCoder of a GroupByKey must be deterministic
核心问题原因
你遗漏了显式标记编码器为确定性的步骤:CustomCoder的默认isDeterministic()方法返回false,即使你的编码逻辑是固定顺序的,Beam依然会判定该编码器是非确定性的,从而拒绝在CoGroupByKey中使用。
此外代码中还有两个次要错误需要修复:
- 变量名重复(第二个PCollection错误复用了
Foo1WithKeys) SimpleFunction的apply方法返回类型错误(错误返回自身类而非KeysId1Id2)
修复方案
1. 修复KeysId1Id2Coder,添加isDeterministic()方法
重写isDeterministic()并返回true,明确告诉Beam该编码器是确定性的:
public class KeysId1Id2Coder extends CustomCoder<KeysId1Id2> { private static final StringUtf8Coder STRING_CODER = StringUtf8Coder.of(); public static KeysId1Id2Coder of() { return new KeysId1Id2Coder(); } @Override public void encode(KeysId1Id2 value, OutputStream outStream) throws IOException { STRING_CODER.encode(value.getId1(), outStream); STRING_CODER.encode(value.getId2(), outStream); } @Override public KeysId1Id2 decode(InputStream inStream) throws IOException { String id1 = STRING_CODER.decode(inStream); String id2 = STRING_CODER.decode(inStream); return new KeysId1Id2(id1, id2); } // 关键:显式标记编码器为确定性 @Override public boolean isDeterministic() { return true; } }
2. 修复主方法中的变量名错误
修正重复的变量名,确保两个PCollection的变量名唯一:
public static PCollection<Output> joinById1Id2(PCollection<Foo1> foo1, PCollection<Foo2> foo2){ PCollection<KV<KeysId1Id2, Foo1>> foo1WithKeys = foo1.apply(WithKeys.of(new MapFoo1ByKeysId1Id2())) .setCoder(KvCoder.of(KeysId1Id2Coder.of(), SerializableCoder.of(Foo1.class))); // 修复变量名重复问题 PCollection<KV<KeysId1Id2, Foo2>> foo2WithKeys = foo2.apply(WithKeys.of(new MapFoo2ByKeysId1Id2())) .setCoder(KvCoder.of(KeysId1Id2Coder.of(), SerializableCoder.of(Foo2.class))); TupleTag<Foo1> foo1TupleTag = new TupleTag<>(); TupleTag<Foo2> foo2TupleTag = new TupleTag<>(); return KeyedPCollectionTuple .of(foo1TupleTag, foo1WithKeys) .and(foo2TupleTag, foo2WithKeys) .apply(CoGroupByKey.create()) .apply(ParDo.of(new JoinFoo1AndFoo2(foo1TupleTag, foo2TupleTag))); }
3. 修复SimpleFunction的返回类型错误
确保apply方法返回的是KeysId1Id2类型,而非自身类:
public class MapFoo1ByKeysId1Id2 extends SimpleFunction<Foo1, KeysId1Id2> { @Override public KeysId1Id2 apply(Foo1 line) { return new KeysId1Id2(line.getId1(), line.getId2()); } } public class MapFoo2ByKeysId1Id2 extends SimpleFunction<Foo2, KeysId1Id2> { @Override public KeysId1Id2 apply(Foo2 line) { return new KeysId1Id2(line.getId1(), line.getId2()); } }
关键说明
Beam要求GroupByKey/CoGroupByKey的KeyCoder必须满足:相同的键对象,每次编码得到的字节序列完全一致。你的编码逻辑已经符合这个要求,但CustomCoder默认不对外声明这一点,必须通过重写isDeterministic()方法明确告知Beam,才能通过校验。
内容的提问来源于stack exchange,提问作者Alberto Martin
相关产品推荐
相关产品推荐

