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

使用kafka-avro-console-producer生产消息失败求助

问题:使用kafka-avro-console-producer导入Avro二进制文件报错

操作流程

  1. 用avro-tools生成Avro二进制文件:
avro-tools fromjson --schema-file foobar-schema.avsc foobar-input.json > foobar.avro
  1. 执行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'"错误。

可行方案

  1. 直接发送原始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"}'即可。

  2. 发送纯Avro二进制文件(不推荐,需额外处理)
    若必须使用已生成的Avro二进制文件,需改用普通的kafka-console-producer,并指定Confluent Avro序列化器的相关配置,同时确保你的Avro二进制文件符合Confluent的消息格式(前4字节为magic byte + Schema ID)。这种方式需要手动处理schema注册和消息格式转换,操作复杂度较高,不如直接发送JSON简便。

内容的提问来源于stack exchange,提问作者Insanovation

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 18:40:02