使用Apache NiFi读取Kafka Topic中Avro数据遇问题求助
Apache NiFi读取Kafka Avro数据报错:AvroRuntimeException: Malformed data length is negative
在使用Apache NiFi从Kafka Topic读取Avro格式数据时,即便使用ConvertAvroToJSON处理器并指定了Schema,仍触发AvroRuntimeException: Malformed data length is negative错误。以下是针对性的排查和解决方法:
检查Kafka数据的序列化方式
如果Kafka中的Avro数据是通过普通Avro序列化写入(而非Confluent Schema Registry的序列化方式),直接用ConvertAvroToJSON会解析失败。建议改用ConsumeKafkaRecord_2_0搭配AvroReader处理:- 配置
ConsumeKafkaRecord_2_0,将Record Reader属性设为AvroReader - 在
AvroReader的配置中填入正确的Avro Schema(支持本地文件路径或内嵌文本) - 后续用
ConvertRecord处理器将Avro格式转成JSON,无需额外指定Schema
- 配置
验证Schema与数据的匹配度
指定的Schema如果和Kafka中Avro数据的结构(字段名、类型、嵌套层级)不匹配,会导致解析时计算数据长度异常:- 核对Schema的每一项定义,确保和生产端的Avro结构完全一致
- 可以用
avro-tools工具验证:先从Kafka导出一条消息到本地文件,执行avro-tools tojson <本地数据文件>,看是否能正常解析为JSON,以此确认Schema的正确性
排查Topic中的异常消息
Kafka Topic可能混入了非Avro格式的消息,或者部分消息在传输中损坏:- 先用
ConsumeKafka_2_0获取原始消息,通过LogAttribute处理器查看消息内容,定位异常消息 - 配置
RouteOnAttribute处理器,添加过滤规则(比如根据消息头、长度等),将异常消息路由到单独的分支处理,避免影响正常数据流程
- 先用
内容的提问来源于stack exchange,提问作者Thiyagaraj Narayanan
相关产品推荐
相关产品推荐

