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

