PostgreSQL 如何同步数据到 Elasticsearch?
PostgreSQL 同步到 Elasticsearch 的方案:Logstash JDBC 定时轮询、逻辑解码 + Debezium 实时 CDC、pgwatch 等工具、或应用双写。本文对比方案并给出 Logstash 配置。
两条主路:定时轮询用 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。