Elasticsearch 如何与 MySQL 集成(数据同步)?
MySQL 数据同步到 Elasticsearch 的四种方案:Logstash JDBC 插件定时拉取、监听 binlog(Canal/Debezium)实时同步、应用双写、以及 CDC 平台。本文对比各方案并给出 Logstash 配置示例。
主流方案是两类:定时轮询——Logstash 的 jdbc input 按更新时间增量拉取;实时同步——Debezium/Canal 监听 binlog 流式写入 ES。小表低频用前者,大表实时性要求高用后者。
方案一:Logstash JDBC 定时拉取
input {
jdbc {
jdbc_connection_string => "jdbc:mysql://mysql:3306/shop"
jdbc_user => "reader"
jdbc_password => "secret"
jdbc_driver_library => "/opt/mysql-connector-j.jar"
statement => "SELECT * FROM products WHERE updated_at > :sql_last_value"
schedule => "*/1 * * * *"
tracking_column => "updated_at"
tracking_column_type => "timestamp"
}
}
output {
elasticsearch {
hosts => ["es:9200"]
index => "products"
document_id => "%{id}" # 用主键当文档 ID,重复同步变更新
}
}
关键技巧:sql_last_value 记录上次同步位点实现增量;document_id 用主键让重复写入变成更新而不是重复文档。
方案二:binlog 实时同步
Debezium(Kafka Connect)或 Canal 监听 MySQL binlog,行变更实时推送到 ES。延迟秒级,但架构复杂度高(多一套 CDC 组件)。
方案三:应用双写
业务代码写完 MySQL 顺手写 ES。简单但耦合——ES 挂了影响主流程,一致性靠重试保障。
选型建议
| 场景 | 方案 |
|---|---|
| 报表/搜索,分钟级延迟可接受 | Logstash JDBC |
| 搜索要求准实时 | binlog CDC |
| 极小项目 | 应用双写 |
常见问题(FAQ)
Q:删除的数据怎么同步? 轮询方案感知不到物理删除——业务上改软删除(加 deleted 标记),或 binlog 方案天然捕获 DELETE。
Q:全量初始化怎么做? 第一次把 tracking 条件去掉全量导一遍,再切增量。