使用Logstash的Kafka输入与Avro Codec解析AVRO失败求助
AVRO本地解析失败的排查与解决
问题背景
尝试跳过Schema Registry,直接用Logstash的avro codec本地解析Kafka中的AVRO数据,但解析后得到的document字段里,id是0、name为空,完全不符合预期。
相关信息
Kafka原始输入数据(JSON形式)
{ "id": 5839355690199358000, "name": "VOH*a|[RF6y?" }
Logstash输出结果
"@timestamp" => 2022-10-08T14:16:12.907374Z, "event" => { "original" => "\u0000\u0000\u0000\u0000\u0002������lj�\u0001\u0018VOH*a|[RF6y?" }, "@version" => "1", "document" => { "id" => 0, "name" => "" } }
Logstash配置
input { kafka { bootstrap_servers => "localhost:9092" topics => ["avro-schema-codec"] codec => avro { schema_uri => "/Users/Smit/Desktop/logstash-8.4.3/config/avro_schema.avsc" target => "[document]" } } } filter { } output { stdout { } }
本地AVSC Schema文件
{ "type": "record", "name": "MyRecord", "namespace": "com.mycompany", "fields": [ { "name": "id", "type": "long" }, { "name": "name", "type": "string" } ] }
核心问题与解决方向
- 数据格式不匹配:看
event.original开头的\u0000\u0000\u0000\u0000\u0002,这是Confluent Schema Registry序列化时加的前缀(前4字节是Magic Byte+Schema ID)。但Logstash原生的avro codec只认纯AVRO二进制数据,处理不了带这个前缀的格式,这就是解析失败的根本原因。 - 换用适配Confluent格式的插件:如果Kafka里的数据是用Confluent的AvroSerializer生成的,就别用原生avro codec了,改用
logstash-codec-avro_confluent插件,这个插件专门处理带Schema Registry前缀的AVRO数据,配置时可以指定本地schema路径,不用连Registry。 - 校验Schema一致性:确保本地AVSC文件的Schema和数据序列化时用的Schema完全一致,字段名、类型、顺序、命名空间都不能错,哪怕差一点都会导致解析出空值或默认值。
- 检查文件权限与路径:确认
schema_uri写的路径是对的,Logstash进程有读这个文件的权限,不然加载不了Schema也会解析失败。
内容的提问来源于stack exchange,提问作者Smit
相关产品推荐
相关产品推荐

