使用kafka-avro-console-producer生产消息失败求助
问题:使用kafka-avro-console-producer导入Avro二进制文件报错
操作流程
- 用avro-tools生成Avro二进制文件:
avro-tools fromjson --schema-file foobar-schema.avsc foobar-input.json > foobar.avro
- 执行kafka-avro-console-producer命令发送该文件:
kafka-avro-console-producer \ --broker-list "foobar" \ --topic "foobar" \ --property schema.registry.basic.auth.credentials.source="USER_INFO" \ --property schema.registry.basic.auth.user.info="foobar:foobar" \ --property schema.registry.url="foobar" \ --property key.schema='{\"type\":\"string\"}' \ --property key.separator=":" \ --property value.schema.file=foobar-schema.avsc \ --property value.schema.registry.url="foobar" \ --property parse.key=true \ < foobar.avro
错误信息
org.apache.kafka.common.errors.SerializationException: Error deserializing json Objavro.schema�{\"type\" to Avro of schema \"string\" at io.confluent.kafka.formatter.AvroMessageReader.readFrom(AvroMessageReader.java:127) at io.confluent.kafka.formatter.SchemaMessageReader.readMessage(SchemaMessageReader.java:397) at kafka.tools.ConsoleProducer$.main(ConsoleProducer.scala:50) at kafka.tools.ConsoleProducer.main(ConsoleProducer.scala) Caused by: com.fasterxml.jackson.core.JsonParseException: Unrecognized token 'Objavro': was expecting (JSON String, Number, Array, Object or token 'null', 'true' or 'false') at [Source: (String)"Obj\u0001\u0004\u0016avro.schema�\u0004{\"type\""; line: 1, column: 11] at com.fasterxml.jackson.core.JsonParser._constructError(JsonParser.java:2418) at com.fasterxml.jackson.core.base.ParserMinimalBase._reportError(ParserMinimalBase.java:759) at com.fasterxml.jackson.core.json.ReaderBasedJsonParser._reportInvalidToken(ReaderBasedJsonParser.java:3038) at com.fasterxml.jackson.core.json.ReaderBasedJsonParser._handleOddValue(ReaderBasedJsonParser.java:2079) at com.fasterxml.jackson.core.json.ReaderBasedJsonParser.nextToken(ReaderBasedJsonParser.java:805) at org.apache.avro.io.JsonDecoder.configure(JsonDecoder.java:124) at org.apache.avro.io.JsonDecoder.<init>(JsonDecoder.java:66) at org.apache.avro.io.JsonDecoder.<init>(JsonDecoder.java:74) at org.apache.avro.io.DecoderFactory.jsonDecoder(DecoderFactory.java:266) at io.confluent.kafka.schemaregistry.avro.AvroSchemaUtils.toObject(AvroSchemaUtils.java:281) at io.confluent.kafka.schemaregistry.avro.AvroSchemaUtils.toObject(AvroSchemaUtils.java:274) at io.confluent.kafka.formatter.AvroMessageReader.readFrom(AvroMessageReader.java:120) ... 3 more
相关文件内容
foobar-schema.avsc
{ "namespace": "com.some.namespace", "type": "record", "name": "SomeEvent", "fields": [ { "name": "fieldOne", "type": "string" }, { "name": "fieldTwo", "type": "double" }, { "name": "fieldThree", "type": ["null", "string"], "default": null } ] }
foobar-input.json
{"fieldOne":"987654321","fieldTwo": 250.75,"fieldThree":null}
已验证操作
已将生成的foobar.avro二进制文件转换回JSON格式,确认数据结构有效。
解决建议
问题根源
kafka-avro-console-producer默认从标准输入读取JSON格式的数据,它会将输入内容按JSON规则解析后,再序列化为带Schema Registry标识的Avro二进制消息发送到Kafka。你传入的foobar.avro是纯Avro二进制文件,工具误将其当作JSON解析,因此抛出"Unrecognized token 'Objavro'"错误。
可行方案
直接发送原始JSON(推荐)
跳过avro-tools生成二进制文件的步骤,直接用foobar-input.json作为kafka-avro-console-producer的输入,工具会自动完成JSON到Confluent格式Avro的转换:kafka-avro-console-producer \ --broker-list "foobar" \ --topic "foobar" \ --property schema.registry.basic.auth.credentials.source="USER_INFO" \ --property schema.registry.basic.auth.user.info="foobar:foobar" \ --property schema.registry.url="foobar" \ --property key.schema='{"type":"string"}' \ --property key.separator=":" \ --property value.schema.file=foobar-schema.avsc \ --property parse.key=true \ < foobar-input.json注意:
key.schema参数用单引号包裹时,内部双引号无需转义,简化为'{"type":"string"}'即可。发送纯Avro二进制文件(不推荐,需额外处理)
若必须使用已生成的Avro二进制文件,需改用普通的kafka-console-producer,并指定Confluent Avro序列化器的相关配置,同时确保你的Avro二进制文件符合Confluent的消息格式(前4字节为magic byte + Schema ID)。这种方式需要手动处理schema注册和消息格式转换,操作复杂度较高,不如直接发送JSON简便。
内容的提问来源于stack exchange,提问作者Insanovation
相关产品推荐
相关产品推荐

