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

