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

Kafka Connect JSONConverter处理Schema Registry返回object类型报错咨询

问题解答

1. 获取到的Schema是否可用于JSONConverter?

你从Schema Registry拿到的是标准JSON Schema规范的结构,无法直接用于Kafka Connect自带的org.apache.kafka.connect.json.JSONConverter,这不是使用方式错误,是二者的Schema定义规范不兼容:

你遇到的报错org.apache.kafka.connect.errors.DataException: Unknown schema type: object正是因为JSONConverter无法识别JSON Schema里的object类型定义。

2. JSONConverter要求的Schema格式是什么?

JSONConverter只识别Kafka Connect内部定义的Schema模型,和标准JSON Schema的差异如下:

  • 类型取值只能是Kafka Connect预定义的枚举值:复合类型用struct(对应JSON Schema的object),基础类型用int32/string/boolean等(对应JSON Schema的integer/string/boolean等)
  • 复合类型的字段定义放在fields数组中,每个字段包含name、type、optional三个核心属性,不使用JSON Schema的properties和required数组配置
  • 必选字段对应optional: false,可选字段对应optional: true

你给出的JSON Schema转换为JSONConverter支持的格式示例如下:

{
  "type": "struct",
  "name": "test-schema",
  "fields": [
    {
      "name": "id",
      "type": "int32",
      "optional": false
    }
  ]
}

3. Schema Registry响应转可用Schema对象的方案

有三种成熟的实现方案:

  • 直接替换转换器:使用Confluent官方提供的io.confluent.connect.json.JsonSchemaConverter,该转换器原生兼容Schema Registry返回的标准JSON Schema,可自动完成到Kafka Connect内部Schema对象的转换,只需要在Connect配置中指定schema.registry.url为你的Schema Registry地址即可
  • 手动编写映射逻辑:按上面提到的格式差异,逐字段把JSON Schema的配置映射为Kafka Connect的Schema结构,比如把object转STRUCT类型、从required数组提取字段的必填属性
  • 调用官方转换工具:引入Confluent的kafka-schema-registry-client依赖,直接用库中内置的JSON Schema到Connect Schema的转换工具类完成转换,无需手动实现映射逻辑

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 13:57:02