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

如何通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 20:33:13