PostgreSQL 如何同步数据到 Elasticsearch?

PostgreSQL 同步到 Elasticsearch 的方案:Logstash JDBC 定时轮询、逻辑解码 + Debezium 实时 CDC、pgwatch 等工具、或应用双写。本文对比方案并给出 Logstash 配置。

最佳实践
PostgreSQL 如何同步数据到 Elasticsearch?封面

两条主路:定时轮询用 Logstash 的 jdbc input(按更新时间增量拉);准实时用 Debezium 读 PG 的 WAL 逻辑解码流式推送。选哪个看实时性要求。

方案一:Logstash JDBC 轮询

input {
  jdbc {
    jdbc_connection_string => "jdbc:postgresql://pg:5432/shop"
    jdbc_user => "reader"
    jdbc_password => "secret"
    jdbc_driver_library => "/opt/postgresql.jar"
    statement => "SELECT * FROM products WHERE updated_at > :sql_last_value"
    schedule => "*/2 * * * *"
    tracking_column => "updated_at"
    tracking_column_type => "timestamp"
    use_column_value => true
  }
}
output {
  elasticsearch {
    hosts => ["es:9200"]
    index => "products"
    document_id => "%{id}"
  }
}

document_id 用主键,重复同步变为更新;updated_at 追踪列实现增量。

方案二:Debezium 逻辑解码(准实时)

PG 开逻辑复制(wal_level = logical),Debezium PG 连接器读 WAL 变更事件,经 Kafka 写入 ES。秒级延迟,支持删除事件——但要多维护 Kafka + Connect 集群。

方案对比

方案 延迟 删除同步 复杂度
JDBC 轮询 分钟级 ❌(需软删除)
Debezium CDC 秒级
应用双写 实时 需自理 低但耦合

常见问题(FAQ)

Q:没有 updated_at 列怎么办? 加列并建触发器维护,或全量定期重建索引。

Q:大表首次全量同步很慢? 分批(按主键范围)并行跑,导入期调大 ES 的 refresh_interval。

获取专属方案

联系我们

加入社区

微信扫码
加入官方交流群

立即体验

在线开通,按量计费,真正的云服务!

立即开始

选择观测云版本

代码托管平台