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

ElasticsearchSinkConnector忽略JSON Schema实现Kafka动态数据完整写入Elastic

问题描述

我们目前在kafka-connect-cluster中使用ElasticsearchSinkConnector,实现数据从Kafka同步到Elastic,架构中已部署JSON-Schemas和SchemaRegistry。
当前业务场景下数据结构非常动态,因此在JSON-Schema中只能将data字段定义为type: object。我们发送如下格式的消息:

{
   "data": {
      "propOne": "hi",
      "propTwo": "hi"
   },
   "type": "Person"
}

使用的Schema定义如下:

{
   "required": ["data", "type"],
   "properties": {
      "type": { "type": "string" },
      "data": { "type": "object" }
   }
}

运行时遇到问题:ElasticSinkConnector仅能获取包含type字段对应字符串值的connect-struct,而data字段始终为空对象。该现象符合预期,因为基于现有Schema最多只能推断出该结构,但我们需要将完整的原始data对象同步到elasticsearch实例,而非仅同步空对象。
基本数据流如下:
数据流示意图
Elasticsearch Sink Connector支持写入Apache Kafka Topic的无Schema JSON数据,但如果使用org.apache.kafka.connect.json.JsonConverter转换器忽略Schema,会出现subjectid字节无法解析的问题。因此我们需要类似org.apache.kafka.connect.json.JsonConverter的组件,可忽略subjectid和magicbyte,直接将原始JSON数据写入Elastic。
已知替代方案:我们可以为该场景/Topic使用无Schema JSON事件,但当前生产者暂不支持该方案,因此希望规避该实现方式。

解决方案

方案1:修改JSON Schema配置(优先选择,改造成本最低)

你当前的Schema中data字段仅声明了对象类型,Kafka Connect的JSON Schema转换器默认不会序列化未在Schema中显式定义的属性,只需要在data字段的定义中新增additionalProperties: true配置即可让转换器保留所有动态字段,修改后的Schema如下:

{
   "required": ["data", "type"],
   "properties": {
      "type": { "type": "string" },
      "data": { 
        "type": "object",
        "additionalProperties": true
      }
   }
}

该方案不需要修改生产者代码、连接器配置,仅需要在Schema Registry中升级对应Topic的Schema版本即可生效,完全兼容现有生产链路。

方案2:自定义消息转换器

如果受内部Schema管理规则限制无法修改现有Schema,可基于org.apache.kafka.connect.json.JsonConverter二次开发自定义转换器,核心逻辑调整反序列化步骤:

  • 跳过消息前5个字节(1位magic byte + 4位subject id)
  • 直接读取剩余字节反序列化为JSON对象
  • 将完整JSON对象转换为Connect Struct传入连接器
    开发完成后将自定义转换器包放入Kafka Connect的plugin目录,在连接器配置中指定value.converter为自定义转换器的全类名即可。

方案3:使用SMT(单消息转换)组件预处理

如果不想做定制开发,可通过第三方SMT插件组合实现需求,示例配置如下:

transforms=parseRawData
transforms.parseRawData.type=com.github.jcustenborder.kafka.connect.transform.json.FromJson$Value
transforms.parseRawData.schema.infer=true
transforms.parseRawData.field=data

该方案需要提前在Kafka Connect集群安装对应SMT插件,不需要修改代码和Schema,适合快速验证场景。


内容的提问来源于stack exchange,提问作者David Leitner

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 13:15:02