【发布时间】: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