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

配置Debezium:如何让PostgreSQL同步到Elasticsearch仅存实体数据?

解决Debezium同步PostgreSQL到Elasticsearch时去除元数据的问题

Debezium默认会把数据库变更事件包装在包含op、before、after等元数据的结构里,要让Elasticsearch直接存储纯Vehicle实体数据,有两种常用配置方式:

方法一:在Debezium源连接器中使用Unwrap变换(推荐)

通过添加ExtractNewRecordState变换,直接从Debezium的事件结构里提取出after字段的实体数据,作为Kafka消息的value。这样整个消息流里都是干净的实体数据,不仅Elasticsearch可以直接使用,其他消费者也能直接读取。

修改PostgreSQL源连接器的配置,添加以下内容:

# 基础源连接器配置(保留你的原有配置)
connector.class=io.debezium.connector.postgresql.PostgresConnector
database.hostname=postgres
database.port=5432
database.user=postgres
database.password=postgres
database.dbname=your_database
database.server.name=postgres-server
table.include.list=public.vehicle

# 添加Unwrap变换配置
transforms=unwrap
transforms.unwrap.type=io.debezium.transforms.ExtractNewRecordState
# 可选配置:删除墓碑消息(针对DELETE事件),根据业务需求调整
transforms.unwrap.drop.tombstones=true
# 可选配置:将DELETE事件重写为null值,方便Elasticsearch自动删除文档
transforms.unwrap.delete.handling.mode=rewrite

方法二:在Elasticsearch Sink连接器中提取目标字段

如果不想修改源端的消息结构,可以在Sink连接器里配置ExtractField变换,指定从消息的after字段读取数据作为Elasticsearch的文档内容。

修改Elasticsearch Sink连接器的配置:

# 基础Sink连接器配置(保留你的原有配置)
connector.class=io.confluent.connect.elasticsearch.ElasticsearchSinkConnector
connection.url=http://elasticsearch:9200
topics=postgres-server.public.vehicle
index.name=vehicles
key.ignore=true

# 配置提取after字段
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
transforms=extractAfter
transforms.extractAfter.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.extractAfter.field=after

验证配置是否生效

可以用Kafka控制台消费者查看消息内容,确认元数据已被去除:

kafka-console-consumer.sh --bootstrap-server kafka:9092 --topic postgres-server.public.vehicle --from-beginning

如果输出的是纯Vehicle实体的JSON结构,说明配置成功。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 11:56:34