Flink KafkaConsumer反序列化复合Avro Schema报ArrayIndexOutOfBoundsException
错误根因
你遇到的java.lang.ArrayIndexOutOfBoundsException: 负数异常是Avro序列化/反序列化格式不匹配导致的典型问题:
- Kafka中存储的Avro消息是通过Confluent Schema Registry序列化生成的,这类消息默认在头部追加了5字节的前缀(1字节魔数+4字节Schema ID)
- 你使用的原生
AvroDeserializationSchema不会识别处理这个前缀,会直接把前缀字节当成Avro数据内容解析,就会触发数组越界异常,不同Schema对应的ID不同,所以异常里的负数值会发生变化。
解决方案
按照下面步骤调整即可解决:
- 替换反序列化器
将原来的AvroDeserializationSchema.forGeneric替换为Confluent协议适配的ConfluentRegistryAvroDeserializationSchema,示例代码如下:
Schema schema = Schema.parse("{\"type\":\"record\",\"name\":\"Envelope\" ... etc"); // 构造适配Confluent Schema Registry的反序列化器,第二个参数填你实际的Schema Registry访问地址 ConfluentRegistryAvroDeserializationSchema<GenericRecord> deserializer = ConfluentRegistryAvroDeserializationSchema.forGeneric( schema, "http://your-schema-registry-host:port" ); DataStreamSource<GenericRecord> stream = env.addSource(new FlinkKafkaConsumer<>(topic, deserializer, getKafkaCredentials()).setStartFromLatest());
- 调整Maven依赖
注释掉原生flink-avro依赖,引入Confluent适配依赖,保证依赖版本和你的Flink版本(1.12.1)完全一致,依赖配置如下:
<!-- 注释掉原生flink-avro --> <!--<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-avro</artifactId> <version>${flink.version}</version> </dependency>--> <!-- 新增Confluent Avro适配依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-avro-confluent-registry</artifactId> <version>${flink.version}</version> </dependency>
你最后附上的完整pom配置已经符合要求,可以直接复用。
内容的提问来源于stack exchange,提问作者AmerS
相关产品推荐
相关产品推荐

