使用@DefaultSchema时WithKeys无法返回默认Coder的问题咨询
问题原因
Beam的Schema自动序列化机制主要针对直接作为PCollection元素类型的POJO,当POJO作为KV<K,V>中的Key时,CoderRegistry在推断KV类型的Coder时,无法自动关联到POJO的SchemaCoder——因为KV是参数化类型,需要分别为Key和Value明确提供Coder,而默认的Coder推断逻辑没有覆盖这种嵌套场景。
解决方案
显式为生成的PCollection<KV<Pojo, String>>指定KvCoder,其中Key的Coder使用SchemaCoder.of(Pojo.class),Value使用对应类型的Coder(这里是StringUtf8Coder)。
修改后的代码示例:
// 处理a1的部分 PCollection<KV<Pojo, String>> a1 = a.apply( WithKeys.<Pojo, String>of(v -> new Pojo(1, v)) .withKeyType(TypeDescriptor.of(Pojo.class))) // 显式指定KV的Coder .setCoder(KvCoder.of(SchemaCoder.of(Pojo.class), StringUtf8Coder.of())); // 处理b1的部分同理 PCollection<KV<Pojo, String>> b1 = b.apply( WithKeys.<Pojo, String>of(v -> new Pojo(1, v)) .withKeyType(TypeDescriptor.of(Pojo.class))) .setCoder(KvCoder.of(SchemaCoder.of(Pojo.class), StringUtf8Coder.of()));
补充说明
你的Pojo类已经正确实现了equals和hashCode方法,这是作为KV Key的必要条件(保证分组/聚合逻辑正确),无需额外修改。
内容的提问来源于stack exchange,提问作者RonVe
相关产品推荐
相关产品推荐

