【问题标题】:Kafka AVRO Deserialization Error when pushed data using Python's Avro library使用 Python 的 Avro 库推送数据时出现 Kafka AVRO 反序列化错误
【发布时间】:2020-10-16 17:54:08
【问题描述】:

我已经设置了一个带有 Kafka 连接节点的 Kafka 集群,该节点具有 Postgres 的接收器配置。

AVRO 架构:

{
    "namespace": "example.avro",
    "type": "record",
    "name": "topicname",
    "fields": [
        {"name": "deviceid", "type": "string"},
        {"name": "longitude", "type": "float"},
        {"name": "latitude",  "type": "float"}
    ]
}

我发布数据的 Python 代码是:

# Path to user.avsc avro schema
SCHEMA_PATH = "user.avsc"
SCHEMA = avro.schema.parse(open(SCHEMA_PATH).read())

writer = DatumWriter(SCHEMA)
bytes_writer = io.BytesIO()
encoder = avro.io.BinaryEncoder(bytes_writer)
writer.write({"deviceid":"9098", "latitude":  90.34 , "longitude": 334.4}, encoder)
raw_bytes = bytes_writer.getvalue()
PRODUCER.send_messages(TOPIC, raw_bytes)

我在 Kafka Connect 日志中收到以下错误:

org.apache.kafka.common.errors.SerializationException: 错误
反序列化 id -1 的 Avro 消息\n原因: org.apache.kafka.common.errors.SerializationException:未知的魔法 字节!\n","id":0,"worker_id":"0.0.0.0:8083"}],"type":"sink"}

可能是什么问题?
或者对于提到的 json 数据,正确的 avro 方案应该是什么?

【问题讨论】:

    标签: java apache-kafka kafka-producer-api


    【解决方案1】:

    我对各种 python 客户端做的不多,但是几乎可以肯定的是,魔术字节错误是因为您发送的可能是有效的 avro,但是如果您想与模式注册表集成,则有效负载需要位于不同的格式(附加标头信息,记录在此处https://docs.confluent.io/current/schema-registry/docs/serializer-formatter.html 搜索有线格式或魔术字节)。我个人会尝试使用 confluent 的 python kafka 客户端——https://github.com/confluentinc/confluent-kafka-python——它有使用 Avro 和模式注册表的示例。

    【讨论】:

      猜你喜欢
      • 2015-08-01
      • 2019-09-11
      • 1970-01-01
      • 2019-12-04
      • 1970-01-01
      • 2019-07-30
      • 2020-06-26
      • 2019-11-18
      相关资源
      最近更新 更多