如何通过Logstash监听MySQL数据变更(无需调度器/手动执行)
无需定时调度:Logstash实时监听MySQL变更的实现方案
你当前的JDBC输入配置基于增量ID轮询实现数据同步,若要实现实时捕获MySQL数据变更、无需手动/定时触发的同步,可采用以下两种成熟方案:
方案一:使用Logstash mysql_binlog 输入插件
该插件直接监听MySQL的二进制日志(Binlog),能实时捕获数据库的INSERT、UPDATE、DELETE操作,完全无需依赖定时调度。
前置准备
开启MySQL Binlog:
修改MySQL配置文件(如my.cnf/my.ini),添加以下配置:server-id = 1 log_bin = /var/log/mysql/mysql-bin.log binlog_format = ROW binlog_row_image = FULL重启MySQL服务后,执行
SHOW VARIABLES LIKE 'log_bin';确认Binlog已启用。授予Logstash账号Binlog读取权限:
GRANT REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'root'@'localhost'; FLUSH PRIVILEGES;
配置示例
替换你当前配置中的input段即可,原有的filter和output可直接保留:
input { mysql_binlog { host => "127.0.0.1" port => 3306 user => "root" password => "pass" database => "db" table => "table" server_id => 100 # 需与MySQL的server-id不同,且在集群中唯一 codec => "json" } }
关键注意事项
binlog_format必须设置为ROW,否则插件无法解析行级变更细节。- 插件会自动记录当前读取的Binlog位置,重启Logstash后将从断点继续同步,无需额外配置追踪字段。
方案二:结合Debezium + Kafka 实现CDC(变更数据捕获)
如果需要更可靠的分布式同步能力,可使用Debezium作为CDC工具捕获MySQL Binlog,将变更事件发送到Kafka,再由Logstash从Kafka消费数据同步至Elasticsearch。该方案适合大规模、高可用的生产场景。
核心流程
- Debezium作为Kafka Connect源连接器,监听MySQL Binlog,将行级变更转为结构化JSON事件发送到Kafka Topic。
- Logstash通过
kafka输入插件订阅对应Topic,处理后同步到Elasticsearch。
Logstash配置示例(仅输入段)
input { kafka { bootstrap_servers => "localhost:9092" topics => "db.table" # Debezium生成的Topic名称,格式为<数据库名>.<表名> codec => "json" } }
方案对比
| 方案 | 优势 | 适用场景 |
|---|---|---|
| mysql_binlog插件 | 配置简单、轻量级,无需额外中间件 | 中小规模、单节点同步场景 |
| Debezium + Kafka | 高可用、消息持久化,支持多下游消费 | 大规模分布式系统场景 |
内容的提问来源于stack exchange,提问作者Mohammad Mudallal
相关产品推荐
相关产品推荐

