You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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相关依赖。

分析与解决方案

核心原因

  1. Avro版本冲突:flink-avro 1.14.2默认依赖Avro 1.10.2,会覆盖你显式声明的1.7.7版本。导致Flink应用实际用高版本Avro序列化数据,而外部应用用1.7.7反序列化,两者二进制格式不兼容,直接触发数组越界异常。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 14:20:23