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
相关产品推荐
相关产品推荐

