如何通过Kafka REST Proxy发送含Union类型的Avro消息?
问题
在Docker容器中使用Kafka、SchemaRegistry、KafkaUI和Kafka REST Proxy做测试时,遇到了Avro Union类型字段的序列化问题。定义的Schema包含Union类型字段,当该字段传null时请求成功,但传递具体数据时两种构造方式都返回400错误,求正确的Union类型数据传递方式。
Avro Schema
{ "type": "record", "name": "MessageValue", "namespace": "com.miwoe.schemas.sandbox", "doc": "Some message", "fields": [ { "name": "message", "type": "string", "doc": "Some message." }, { "name": "unionTypeField", "type": [ "null", { "type": "record", "name": "UnionTypeField", "fields": [ { "name": "someField", "type": "string" } ] } ], "default": null } ] }
成功的请求体(Union字段为null)
{ "value_schema_id": 7, "records": [ { "value": { "message": "foo", "unionTypeField": null } } ] }
失败的请求体1及错误
{ "value_schema_id": 7, "records": [ { "value": { "message": "foo", "unionTypeField": { "UnionTypeField" : { "someField": "foo" } } } } ] }
返回错误:
{ "error_code": 400, "message": "Bad Request: Unknown union branch UnionTypeField" }
失败的请求体2及错误
{ "value_schema_id": 7, "records": [ { "value": { "message": "foo", "unionTypeField": { "someField": "foo" } } } ] }
返回错误:
{ "error_code": 400, "message": "Bad Request: Unknown union branch someField" }
解决方案
问题出在Avro Union类型的JSON序列化规则上:当Union类型的分支是自定义record时,必须使用全限定类型名(namespace + 类型名)作为键来包裹具体的record数据,而不是仅用类型名或直接传record内容。
正确的请求体应该这样构造:
{ "value_schema_id": 7, "records": [ { "value": { "message": "foo", "unionTypeField": { "com.miwoe.schemas.sandbox.UnionTypeField": { "someField": "foo" } } } } ] }
关键说明
- Avro的Union类型在JSON序列化时,除了
null可以直接写值,其他分支必须用对应的类型标识来包裹。对于自定义record,这个标识是全限定名(即schema中定义的namespace加上name)。 - 之前失败的请求体1只用了
UnionTypeField(缺少namespace),请求体2直接传了record内容,都不符合Avro的JSON序列化规范,所以被Kafka REST Proxy判定为无效请求。
内容的提问来源于stack exchange,提问作者miwoe
相关产品推荐
相关产品推荐

