You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过Logstash监听MySQL数据变更(无需调度器/手动执行)

无需定时调度:Logstash实时监听MySQL变更的实现方案

你当前的JDBC输入配置基于增量ID轮询实现数据同步,若要实现实时捕获MySQL数据变更、无需手动/定时触发的同步,可采用以下两种成熟方案:


方案一:使用Logstash mysql_binlog 输入插件

该插件直接监听MySQL的二进制日志(Binlog),能实时捕获数据库的INSERT、UPDATE、DELETE操作,完全无需依赖定时调度。

前置准备

  1. 开启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已启用。

  2. 授予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。该方案适合大规模、高可用的生产场景。

核心流程

  1. Debezium作为Kafka Connect源连接器,监听MySQL Binlog,将行级变更转为结构化JSON事件发送到Kafka Topic。
  2. 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.22 21:34:57