Debezium同步PostgreSQL JSONB到Elasticsearch嵌套对象问题求助
处理PostgreSQL JSONB字段到Elasticsearch嵌套对象的CDC同步问题
问题场景
在K8s集群通过Strimzi部署CDC同步链路:使用io.debezium.connector.postgresql.PostgresConnector作为源连接器捕获PostgreSQL数据,同步到Kafka Topic后,再通过io.confluent.connect.elasticsearch.ElasticsearchSinkConnector写入Elasticsearch。核心问题是PostgreSQL中的JSONB字段经同步后在Kafka中以字符串形式存在,无法直接映射为Elasticsearch的嵌套对象。
当前现象
- Kafka Topic中的实际消息格式:
{ "id": "someId", "title": "someTitle", "someJsonBField": "[{\"aField\": \"test\"},{\"aField\": \"test\"}]" }
- 期望的消息格式:
{ "id": "someId", "title": "someTitle", "someJsonBField": [ { "aField": "test" }, { "aField": "test" } ] }
- 预映射索引时的报错:
Indexing failed: ElasticsearchException[Elasticsearch exception [type=document_parsing_exception, reason=[1:16] object mapping for [someJsonBField] tried to parse field [someJsonBField] as object, but found a concrete value]]
- 动态映射的问题:Elasticsearch会将JSONB字段识别为字符串类型,不符合嵌套对象的需求。
现有连接器配置
源连接器(Debezium PostgreSQL)
spec: class: io.debezium.connector.postgresql.PostgresConnector tasksMax: 2 config: topic.prefix: xxx table.include.list: xxx database.hostname: xxx database.port: 5432 database.user: xxx database.password: xxx database.dbname: dbName slot.name: replication_slot publication.name: publication_name decimal.handling.mode: double plugin.name: pgoutput snapshot.mode: initial value.converter: "org.apache.kafka.connect.json.JsonConverter" value.converter.schemas.enable: "false" transforms: unwrap transforms.unwrap.type: "io.debezium.transforms.ExtractNewRecordState" transforms.unwrap.drop.tombstones: "false"
Sink连接器(Elasticsearch)
spec: class: io.confluent.connect.elasticsearch.ElasticsearchSinkConnector tasksMax: 2 config: connector.class: "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector" tasks.max: "2" topics: xxx connection.url: xxx connection.username: xxx connection.password: xxx flush.synchronously: "true" value.converter: "org.apache.kafka.connect.json.JsonConverter" value.converter.schemas.enable: "false" key.ignore: "false" schema.ignore: "true" errors.tolerance: "all" errors.deadletterqueue.topic.name: "failed-records-dlq" errors.deadletterqueue.context.headers.enable: "true" errors.deadletterqueue.topic.replication.factor: "3" errors.log.enable: "true" errors.log.include.messages: "true"
解决方案:使用Kafka Connect内置JSON解析SMT
无需自定义SMT,利用Kafka Connect自带的Json$Value转换器即可将JSONB字符串字段解析为嵌套JSON结构。
修改源连接器配置
在原有transforms基础上添加JSON解析的转换规则:
transforms: unwrap,parseJson # 保留原有的unwrap配置 transforms.unwrap.type: "io.debezium.transforms.ExtractNewRecordState" transforms.unwrap.drop.tombstones: "false" # 新增JSON解析转换 transforms.parseJson.type: org.apache.kafka.connect.transforms.Json$Value transforms.parseJson.fields: someJsonBField # 替换为你的JSONB字段名,多个字段用逗号分隔(如field1,field2)
说明
- 该SMT会将指定字段的JSON字符串解析为对应的JSON对象/数组,确保Kafka消息中的JSONB字段以结构化形式存在。
- 若存在多个JSONB字段,只需在
fields参数中用逗号分隔字段名即可。 - 修改配置后,Kafka消息格式将符合预期,Elasticsearch的预映射可正确识别为嵌套对象,解决解析报错问题。
内容的提问来源于stack exchange,提问作者Septiyan
相关产品推荐
相关产品推荐

