Apache Beam自定义DestinationT类编码非确定性异常求助
这个问题我之前碰到过好几次,核心原因是Apache Beam对所有编码器(Coder)有个硬性要求:必须是确定性的——简单说就是同一个对象每次编码出来的字节序列得完全一模一样,不然像窗口计算、键值分组、容错恢复这些依赖数据一致性的操作都会出问题。而Java原生的SerializableCoder因为序列化过程受类结构、字段顺序甚至JVM版本影响,没法保证这种一致性,所以才抛出了这个异常。
下面给你几个靠谱的解决办法,按适配性排序(毕竟你在用AvroIO,优先选和Avro适配的方案):
方法一:用AvroCoder适配自定义Bean类(首推)
既然你已经在使用AvroIO,直接用Avro的序列化机制来处理你的DestinationT Bean是最顺理成章的,步骤也很清晰:
给自定义Bean加Avro注解:
给你的DestinationT类添加上Avro的注解,让Avro能识别它的结构,注意必须要有无参构造器(Avro序列化需要):import org.apache.avro.reflect.Nullable; import org.apache.avro.reflect.Record; @Record public class DestinationT implements Serializable { @Nullable // 标记允许为空的字段 private String eventType; private Long timestamp; // 其他业务字段... // 必须要有无参构造器,Avro序列化会用到 public DestinationT() {} // 标准的Getter和Setter方法 public String getEventType() { return eventType; } public void setEventType(String eventType) { this.eventType = eventType; } public Long getTimestamp() { return timestamp; } public void setTimestamp(Long timestamp) { this.timestamp = timestamp; } }注册AvroCoder到你的Pipeline:
你可以全局注册这个Coder,或者在具体的PCollection上指定:// 全局注册,整个Pipeline都能用 PipelineOptions options = PipelineOptionsFactory.create(); CoderRegistry coderRegistry = options.getCoderRegistry(); coderRegistry.registerCoderForClass(DestinationT.class, AvroCoder.of(DestinationT.class)); // 或者针对特定的PCollection指定 PCollection<DestinationT> eventCollection = ...; eventCollection.setCoder(AvroCoder.of(DestinationT.class));
这样AvroIO就能用确定性极强的AvroCoder来处理你的自定义Bean,彻底避开SerializableCoder的坑。
方法二:自定义确定性Coder
要是你不想依赖Avro注解,也可以自己写一个符合Beam要求的Coder,手动控制每个字段的编码顺序和方式,确保一致性:
实现Coder接口:
针对DestinationT编写自定义Coder,注意编码和解码的顺序必须完全一致,而且用到的子编码器也得是确定性的:import org.apache.beam.sdk.coders.Coder; import org.apache.beam.sdk.coders.CoderException; import org.apache.beam.sdk.coders.StringUtf8Coder; import org.apache.beam.sdk.coders.VarLongCoder; import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; public class DestinationTCoder extends Coder<DestinationT> { // 用单例模式避免重复创建实例 private static final DestinationTCoder INSTANCE = new DestinationTCoder(); // 复用Beam提供的确定性基础编码器 private static final StringUtf8Coder STRING_CODER = StringUtf8Coder.of(); private static final VarLongCoder LONG_CODER = VarLongCoder.of(); public static DestinationTCoder of() { return INSTANCE; } @Override public void encode(DestinationT value, OutputStream outStream) throws CoderException, IOException { // 按固定顺序编码字段,绝对不能乱序 STRING_CODER.encode(value.getEventType(), outStream); LONG_CODER.encode(value.getTimestamp(), outStream); // 其他字段也按同样的方式依次编码 } @Override public DestinationT decode(InputStream inStream) throws CoderException, IOException { DestinationT dest = new DestinationT(); // 解码顺序必须和encode完全对应 dest.setEventType(STRING_CODER.decode(inStream)); dest.setTimestamp(LONG_CODER.decode(inStream)); // 其他字段按同样顺序解码 return dest; } @Override public void verifyDeterministic() throws NonDeterministicException { // 验证所有用到的子编码器都是确定性的 STRING_CODER.verifyDeterministic(); LONG_CODER.verifyDeterministic(); } }注册自定义Coder:
注册方式和AvroCoder一样,在Pipeline中绑定到你的DestinationT类:PipelineOptions options = PipelineOptionsFactory.create(); options.getCoderRegistry().registerCoderForClass(DestinationT.class, DestinationTCoder.of());
方法三:禁用确定性检查(强烈不推荐,仅临时测试用)
如果只是临时跑个测试图省事,你可以通过PipelineOptions关掉编码器的确定性检查,但生产环境绝对别这么做——这会破坏Beam的核心保证,导致数据重复、窗口计算错误、容错失败等各种诡异问题:
PipelineOptions options = PipelineOptionsFactory.create(); options.setEnforceDeterministicCoders(false);
内容的提问来源于stack exchange,提问作者ASe

