Flink序列化的Avro消息外部反序列化失败问题排查
问题
我有一个生成的Avro类PeopleTransformation(继承自SpecificRecordBase),未使用Schema Registry,希望无需它解决问题。在Flink应用中使用以下Avro依赖显式序列化该类实例:
<dependency> <groupId>org.apache.avro</groupId> <artifactId>avro</artifactId> <version>1.7.7</version> </dependency>
序列化代码如下(暂未添加try-catch-finally,属单独问题):
public static byte[] serialize(PeopleTransformation record) throws IOException { DatumWriter<PeopleTransformation> writer = new SpecificDatumWriter<PeopleTransformation>(record.getSchema()); ByteArrayOutputStream out = new ByteArrayOutputStream(); BinaryEncoder encoder = EncoderFactory.get().binaryEncoder(out, null); writer.write(record, encoder); encoder.flush(); IOUtils.closeQuietly(out); return out.toByteArray(); }
将生成的字节数组通过FlinkKafkaProducer发送至Kafka,在外部非Flink Java应用中使用相同Avro库反序列化,代码如下:
public static <T extends SpecificRecordBase> T deserializeRecord(byte[] bytes, Class<T> avroClass) throws IOException { SpecificDatumReader<T> reader = new SpecificDatumReader<>(avroClass); return reader.read(null, DecoderFactory.get().binaryDecoder(bytes, null)); }
操作抛出如下异常:
java.lang.ArrayIndexOutOfBoundsException: 40 at org.apache.avro.io.parsing.Symbol$Alternative.getSymbol(Symbol.java:402) at org.apache.avro.io.ResolvingDecoder.doAction(ResolvingDecoder.java:290) ...... at org.apache.avro.generic.GenericDatumReader.read(GenericDatumReader.java:142)
经排查,Flink应用包含以下依赖(虽未显式使用其类,但不可缺少):
<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-avro</artifactId> <version>1.14.2</version> </dependency>
该包引入了org.apache.flink.formats.avro.typeutils.AvroTypeInfo类,会自动将Schema中的每个CharSequence转换为Utf8。移除该依赖会出现AvroTypeInfo缺失异常,请问该替换是否导致了兼容性问题?或有其他原因?生产环境中消息将被第三方应用消费,难以引入任何Flink相关依赖。
分析与解决方案
核心原因
- Avro版本冲突:
flink-avro 1.14.2默认依赖Avro 1.10.2,会覆盖你显式声明的1.7.7版本。导致Flink应用实际用高版本Avro序列化数据,而外部应用用1.7.7反序列化,两者二进制格式不兼容,直接触发数组越界异常。 - Flink隐式类型转换干扰:
AvroTypeInfo强制将字符串类型转为Utf8,高版本Avro对Utf8的序列化逻辑和1.7.7存在差异,进一步加剧了兼容性问题。
解决步骤
- 强制统一Avro版本:在Flink应用的pom.xml中通过
dependencyManagement锁定Avro版本为1.7.7,避免Flink依赖的高版本Avro覆盖:
<dependencyManagement> <dependencies> <dependency> <groupId>org.apache.avro</groupId> <artifactId>avro</artifactId> <version>1.7.7</version> </dependency> </dependencies> </dependencyManagement>
- 隔离Flink Avro的隐式处理:如果不需要Flink的Avro序列化支持,可尝试排除
flink-avro中的Avro依赖;若必须保留该依赖,需确保序列化过程完全由自定义代码控制——不要让Flink的DataStream类型系统处理PeopleTransformation实例,提前取出对象手动序列化后再发送到Kafka。 - 验证序列化一致性:在Flink应用内序列化数据后,直接用1.7.7版本的Avro库本地反序列化测试,确认数据可正常解析后再发送至Kafka。
内容的提问来源于stack exchange,提问作者Paul O'Neill
相关产品推荐
相关产品推荐

