尧图精选

Logstash实战:MySQL/MariaDB数据同步到Elasticsearch全攻略

🕒 发布时间:2026/9/9 10:53:36 📁 来源:尧图网络
前阵子做电商订单库的检索重构要把 MySQL 里几千万订单数据搬到 Elasticsearch 里做快速搜索和聚合。表结构、同步策略、增量更新的方案选了一圈最后选了 Logstash 作为主力同步工具。如果你也遇到类似的场景——手头数据在 MySQL/MariaDB 里躺着业务方却天天要各种维度搜索和统计——这篇文章应该能在半个小时内帮你跑通一条可上线的同步链路。全文我会按“方案选型→环境准备→第一个同步→增量→字段映射→调优→踩坑”的顺序来讲涉及的命令和配置都会给全你可以直接抄作业。顺便说一句MariaDB 是 MySQL 的分支JDBC 层面略有差异但整体思路完全一致下面遇到差异我会单独标注。1. 先想清楚为什么拿 Logstash 做 MySQL/MariaDB 到 Elasticsearch 的同步1.1 同步需求从哪来很多业务系统把数据放在 MySQL 或 MariaDB 里这没问题但一旦业务方开始提这些需求你就会发现 MySQL 真心扛不住按商品名称、备注、描述做全文搜索用 LIKE %keyword% 在千万级表上扫慢到怀疑人生同时组合十几个筛选条件比如“近30天、金额大于500、地区是华东、包含某关键词”这种查询在 MySQL 里要优化半天在 Elasticsearch 里一个 bool 查询就解决还要按天、按小时、按品类做聚合统计MySQL 的 GROUP BY 出报表数据量大一点就卡死。这些需求本质上是搜索和分析的场景。Elasticsearch 天生为这个设计的倒排索引、分布式聚合、毫秒级响应都比 MySQL 硬扛要舒服得多。所以核心问题变成数据怎么从 MySQL/MariaDB 持续进入 Elasticsearch这就引出了 Logstash。1.2 四种主流方案怎么选同步方案我调研过四种简单对比一下方案实时性引入成本维护成本适合场景Logstash JDBC 轮询分钟级可调到秒级低低中小数据量、准实时同步、搜索索引构建Canal Kafka 消费端秒级/毫秒级高高要求实时性高、需要 binlog 级精确同步DataX / 离线批任务离线按小时/天中中大数据量一次性全量导入自研脚本取决于你低起步高逻辑极度定制化、不想依赖任何框架我最后选了 Logstash原因很直接不想为了同步数据专门维护一套 Canal Kafka 的链路。Logstash 是 ELK 体系自带的和 Elasticsearch 天然兼容JDBC 输入插件能直接连 MySQL/MariaDB配置一个文件就能跑起来。对于“数据不需要秒级实时分钟级刷新完全够用”的业务这是性价比最高的选择。1.3 Logstash 同步的边界用 Logstash 做同步之前有几件事必须提前知道否则后面会被坑它是定时轮询不是 binlog 订阅。JDBC 插件按 schedule 配置的 cron 表达式周期性执行 SQL把查询结果写入 Elasticsearch所以实时性最好也就秒级通常是分钟级。它只能同步“查询得到的数据”删除操作同步不了。MySQL 里物理删掉一行Elasticsearch 里对应文档不会自动删。复杂的数据转换能力有限。虽然 pipeline 里有 filter 可以做 mutate、grok、ruby 处理但你要在里面写复杂的多表关联逻辑会非常痛苦不如在 SQL 里直接处理好。超大表全量导入需要分批策略。几千万行一次性 SELECT数据库内存和 Logstash JVM 都可能被压垮后面我会讲怎么拆。一句话总结Logstash 适合“准实时增量、全量先行、逻辑集中在 SQL 层”的同步场景。如果业务要求秒级实时 删除同步直接选 Canal 方案别在 Logstash 上浪费时间。2. 环境准备把 Logstash 和 JDBC 驱动一次配到位2.1 安装 Logstash 并确认版本兼容性版本兼容性是我每次部署都要先确认的事。Logstash 和 Elasticsearch 的版本必须保持大版本一致最佳实践是同 minor 版本。比如你 ES 是 7.17.xLogstash 也尽量用 7.17.xES 8.x 配 Logstash 8.x。版本不匹配时output 插件调用 ES REST API 的兼容性可能出问题比如 7.x 的 Logstash 写 ES 8 集群ES 默认开启安全认证后连接会被拒。安装本身不复杂。Linux 下最省事的是解压 tar.gzwget https://artifacts.elastic.co/downloads/logstash/logstash-8.13.0-linux-x86_64.tar.gz tar -zxvf logstash-8.13.0-linux-x86_64.tar.gz cd logstash-8.13.0Windows 上直接下载 zip 包解压运行bin\logstash.bat。Logstash 8.x 的发行包自带 JDK不用单独装 Java如果你用的是 7.x系统里最好有 Java 11 或 17。装完跑一下版本号验证bin/logstash --version看到版本输出基本就 OK 了。JDBC 插件logstash-input-jdbc在官方发行包里已经内置不需要额外安装离线环境如果缺插件再单独用bin/logstash-plugin install logstash-input-jdbc补。2.2 JDBC 驱动的选择MySQL 和 MariaDB 不一样很多人第一步就卡在驱动上。Logstash 的 JDBC 插件不负责实现数据库通信协议连接 MySQL 要用 MySQL 官方驱动连接 MariaDB 要用 MariaDB 的 Java Client这俩不能混用。MySQL 这边新项目请直接用 MySQL Connector/J 8.x也就是mysql-connector-j-8.0.33.jar这类文件。驱动的类名是com.mysql.cj.jdbc.Driver连接串前缀是jdbc:mysql://。注意老版本 5.x 驱动的类名是com.mysql.jdbc.Driver没有cj很多老教程这么写配 8.x 驱动时会报 ClassNotFoundException。MariaDB 这边用官方mariadb-java-client-x.x.x.jar驱动类是org.mariadb.jdbc.Driver连接串前缀是jdbc:mariadb://。虽然 MariaDB 官方说大部分情况下能兼容 MySQL 驱动但我实际试下来在认证方式、时间类型处理上都可能出幺蛾子所以别省这一步。下载 jar 后放到一个 Logstash 进程有权限读取的固定目录比如logstash-8.13.0/driver/或者/usr/share/logstash/driver/。注意不要随便扔在 home 目录下Logstash 如果是 systemd 启动可能读不到。2.3 验证驱动能被 Logstash 正常加载先别急着写完整配置写个最小验证配置让 Logstash 能连接数据库并读一行数据确认驱动没问题。直接这样跑bin/logstash -e input { jdbc { jdbc_driver_library /path/to/mysql-connector-j-8.0.33.jar jdbc_driver_class com.mysql.cj.jdbc.Driver jdbc_connection_string jdbc:mysql://127.0.0.1:3306/test?useSSLfalseserverTimezoneAsia/Shanghai jdbc_user root jdbc_password password statement SELECT 1 AS id } } output { stdout { codec rubydebug } } --config.test_and_exit--config.test_and_exit是 Logstash 的配置检查模式只校验配置是否正确不会真正跑数据流。如果这一步能通过说明 jar 包路径、驱动类名、连接串这三个最容易错的地方都对了。如果报错大概率是下面几个原因jar 包路径写错或权限不足驱动类名和驱动 jar 版本不匹配数据库 URL 里的参数有语法问题。这一步排查清楚后面所有配置都能在这个基础上叠加。3. 跑通第一个同步任务MySQL 单表全量导入 ES3.1 最小配置长什么样环境没问题后直接上一个能用的最小同步配置。我以订单表orders为例表结构大致是id、order_no、user_id、status、create_time。先做全量导入不考虑增量input { jdbc { jdbc_driver_library /usr/share/logstash/driver/mysql-connector-j-8.0.33.jar jdbc_driver_class com.mysql.cj.jdbc.Driver jdbc_connection_string jdbc:mysql://127.0.0.1:3306/shop?useSSLfalseserverTimezoneAsia/ShanghaicharacterEncodingutf8 jdbc_user logstash jdbc_password your_password schedule */5 * * * * * statement SELECT id, order_no, user_id, status, create_time FROM orders clean_run true } } output { elasticsearch { hosts [http://127.0.0.1:9200] index orders_index document_id %{id} } }这个配置存成orders.conf然后启动bin/logstash -f orders.conf第一次跑就是把orders表全量灌进orders_index。需要说明的是这个配置看起来简单但已经包含了同步任务最核心的骨架input从数据库读数据output写入 Elasticsearch。后面所有增量、字段处理、多表同步都是在这两个模块上做文章。3.2 每个关键参数的作用我拆几个最核心的参数讲讲为什么要这么配。jdbc_driver_library和jdbc_driver_class前面讲过驱动路径和类名必须一一对应。jdbc_connection_string里我显式加了useSSLfalse和serverTimezoneAsia/Shanghai。这两个参数几乎必加不加 SSLMySQL 8 默认可能报 SSL 连接错误不加 serverTimezone日期字段经常出现 8 小时时差。characterEncodingutf8是防止中文乱码虽然现在 Connector/J 8.x 默认 UTF-8写上更保险。schedule */5 * * * * *是 cron 表达式注意这里是 6 位最后一位是秒。这个表达式的意思是每 5 秒执行一次 SQL。在实际项目里我通常根据业务实时性要求调节比如每 1 分钟同步一次写0 */1 * * * *。statement是查询语句。全量同步时直接SELECT全部字段但真实场景千万别写SELECT *明确列出需要的列既减少网络传输也方便后续维护字段映射。clean_run true的意思是忽略上次同步记录重新从零开始。这里我们是全量导入所以设成 true。等做增量同步时要改成 false让它读取上次同步位置。output.elasticsearch里的hosts是 ES 地址index是索引名document_id是最关键的一行我把 MySQL 表的主键id映射为 ES 文档 ID。这样同步任务重复执行时同一id的文档会覆盖更新不会产生重复数据。如果这一行不写ES 会自动生成随机文档 ID跑两次全量就有两份一模一样的文档。3.3 启动任务并验证数据启动后观察日志重点看有没有[INFO] ... Pipeline started和[INFO] ... Elasticsearch ... 200 OK之类的输出。第一次同步如果表数据量大需要等一会儿。数据验证我一般两个入口命令行用 curl或者 Kibana Dev Tools 执行GET orders_index/_countGET orders_index/_search { query: { match_all: {} }, size: 10 }对比 MySQLSELECT COUNT(1) FROM orders;看两边数量是否一致。第一次全量同步如果 count 对不上优先怀疑document_id设置、字段类型映射和 SQL 条件这几个是最常见的翻车点。4. 增量同步核心机制与实战细节4.1 增量同步的两条路线跑通全量之后业务方就会说“数据要能更新啊不能每次全量导吧。”没错全量同步只适合一次性初始导入日常必须走增量。Logstash JDBC 插件的增量机制核心就一句话用:sql_last_value占位符记住上次查到的位置下次只查这个位置之后的数据。实现方式有两种二选一实现方式配置要点适用场景基于自增主键use_column_value truetracking_column idtracking_column_type numeric只追加、不更新的流水表比如订单日志、操作记录基于更新时间戳use_column_value truetracking_column update_timetracking_column_type timestamp经常修改的业务表比如订单状态、用户资料两种方式的核心配置差异只在tracking_column指哪个字段、tracking_column_type是numeric还是timestamp。我实际项目里绝大多数业务表都优先用update_time字段。因为自增主键只能解决“新增”问题解决不了“更新”问题——行的主键没变但内容变了用 id 跟踪就漏掉了。update_time字段只要记录被修改就会变化配合WHERE update_time :sql_last_value新增和更新都能抓到。4.2 sql_last_value 与 last_run_metadata_path 的行为增量配置的完整示例基于上面订单表加一个update_time字段input { jdbc { jdbc_driver_library /usr/share/logstash/driver/mysql-connector-j-8.0.33.jar jdbc_driver_class com.mysql.cj.jdbc.Driver jdbc_connection_string jdbc:mysql://127.0.0.1:3306/shop?useSSLfalseserverTimezoneAsia/Shanghai jdbc_user logstash jdbc_password your_password schedule */1 * * * * * statement SELECT id, order_no, user_id, status, update_time FROM orders WHERE update_time :sql_last_value use_column_value true tracking_column update_time tracking_column_type timestamp last_run_metadata_path /usr/share/logstash/conf/.orders_last_run clean_run false } } output { elasticsearch { hosts [http://127.0.0.1:9200] index orders_index document_id %{id} } }这里WHERE update_time :sql_last_value中的:sql_last_value是 JDBC 插件自动维护的变量值来自上一次执行结束后记录的位置。位置存在last_run_metadata_path指定的文件里。关于这个机制有几个细节必须说明白第一clean_run false时插件每次启动会读取last_run_metadata_path里的值作为本次sql_last_value的初始值如果文件不存在会从“最早”开始。对 timestamp 类型初始值是1970-01-01 00:00:00左右对 numeric 类型初始值是 0。第二clean_run true会强制忽略已有记录重置为初始值一般只在调试或重新全量时用。第三同一份 Logstash 配置里如果你同时跑多个jdbc input比如订单表和用户表都要同步每个 input 必须指定不同的last_run_metadata_path。否则多个 input 会互相覆盖同步位置文件导致某些表每次都从最开始扫增量完全失效。这是最常见的低级错误。第四边界问题。会导致同一秒内被多次修改的数据重复拉取但因为 ES 按document_id覆盖写入重复拉取并不会产生重复数据最多是多余的网络开销可以接受。4.3 删除数据到底该怎么同步这是 Logstash 方案最明显的短板。JDBC 插件执行 SQL 只能查到“还在表里”的数据MySQL 里物理删除的行ES 里的文档永远不会消失。如果业务对删除同步有强需求我给三个方案软删除这是最推荐的做法。MySQL 表加一个deleted字段删除时执行 UPDATE 而不是 DELETESQL 里写WHERE deleted 0同步时把deleted也写进 ES 文档。搜索时过滤deleted: false。这样 Logstash 能正常增量同步删除行为代价是业务代码需要配合改造。定期重建索引适合索引可以短暂丢失的场景。比如每天凌晨跑一次全量重建白天增量更新。这个方案简单粗暴但数据量大了之后全量重建时间很长不适用。换方案如果必须做到秒级、精确的物理删除同步不要再用 Logstash 硬扛直接上 Canal 监听 binlog下游消费删除事件。两者不是一个层级的东西。我的经验是90% 的业务搜索场景软删除都够用。技术选型一定先看业务容忍度别为了 10% 的需求把整个架构搞复杂。5. 字段处理与索引映射别把同步当成搬家5.1 类型映射MySQL 和 ES 的数据类型对照数据同步不进 ES 不是终点字段类型对不对直接决定搜索和聚合能不能用。ES 的动态映射dynamic mapping虽然能自动识别类型但经常识别的不是你想要的结果。MySQL 和 ES 的类型对照我整理了一份常用表MySQL 类型ES 动态映射结果建议INT / INTEGERinteger计数类字段够用BIGINTlong默认没问题VARCHARtext keywordES 7.x 起默认双字段keyword 用于精确匹配TEXTtext全文搜索DATETIME / TIMESTAMPdate注意时间格式容易被动态映射成 textDECIMAL可能是 double 或 text金额建议用 scaled_float精度可控TINYINT(1)boolean布尔值注意是否为 0/1JSON可能被当字符串用 json filter 解析成对象最常踩的坑是DATETIME。如果 MySQL 里的时间格式不标准或者驱动时区没设对ES 动态映射出来的 date 字段可能格式不对Kibana 里显示为 null 或者日期排序错乱。再比如 DECIMAL 金额字段动态映射成 double 会出现精度尾差做金额聚合时误差虽然小但财务报表不认。所以规范的做法是对重要索引提前写好 mapping别什么都丢给动态映射。5.2 嵌套结构与关联数据的同步前面说的是单表字段实际业务往往涉及多表关联。订单表在 ES 里希望长这样{ id: 1001, order_no: SO20250101001, items: [ { sku: A001, name: 手机, qty: 2 }, { sku: B002, name: 耳机, qty: 1 } ], total_amount: 5999.00 }MySQL 里订单和订单明细是两张表。如果只同步 orders 表明细丢了如果同步明细每个明细变成一条 ES 文档查询订单又得聚合。方案是在 SQL 层把子表组装成 JSON 字符串然后让 Logstash 把 JSON 解析成嵌套对象。MySQL 5.7 和 MariaDB 10.5 都支持 JSON_ARRAYAGG 和 JSON_OBJECTSELECT o.id, o.order_no, o.total_amount, ( SELECT JSON_ARRAYAGG(JSON_OBJECT(sku, i.sku, name, i.name, qty, i.qty)) FROM order_item i WHERE i.order_id o.id ) AS items FROM orders o WHERE o.update_time :sql_last_value重点来了Logstash 拿到这个items字段它只是一个 JSON 字符串不是对象。如果不处理ES 会把items存成 text 或者 keyword没法做嵌套查询。需要在 filter 里加一个json插件把字符串转成结构化对象filter { json { source items target items } }经过这一步ES 里的items才会被正确索引为数组对象才能用 nested 查询或 object 查询。如果 MySQL 版本不支持 JSON_ARRAYAGG那只能手动用 CONCAT 拼接 JSON 字符串或者升级版本总之这步不能省。5.3 自定义索引模板的落地方式为了避免动态映射乱来最稳的方式是提前在 ES 里建好索引模板。我通常是先在 Kibana Dev Tools 里手动 PUT 模板然后让 Logstash 只负责写数据。模板文件大概长这样PUT _index_template/orders_template { index_patterns: [orders_index*], template: { settings: { number_of_shards: 3, number_of_replicas: 1 }, mappings: { properties: { id: { type: long }, order_no: { type: keyword }, user_id: { type: long }, status: { type: keyword }, total_amount: { type: scaled_float, scaling_factor: 100 }, update_time: { type: date, format: yyyy-MM-dd HH:mm:ss||strict_date_optional_time||epoch_millis } } } } }模板建好之后Logstash 配置里可以关掉它的自动模板管理避免和自定义模板冲突output { elasticsearch { hosts [http://127.0.0.1:9200] index orders_index document_id %{id} manage_template false } }这里manage_template false的意思是Logstash 不再自动上传自带的索引模板。如果你不关Logstash 默认会尝试创建一个以logstash命名的模板可能导致字段类型和你预定义的不一致。这个开关在第一次写 index 时就会生效所以一定要在正式同步前确认模板已存在。6. 性能调优与一致性校验6.1 管好 Logstash 的内存和并发同步任务跑久了第一个问题往往是性能跟不上。Logstash 默认 JVM 堆内存只有 1GB数据量大点就容易触发 GC 频繁或 OOM。调内存直接改config/jvm.options-Xms2g -Xmx2g-Xms 和 -Xmx 建议设成一样避免运行期动态扩容引发停顿。具体设多大看你的机器内存和单条数据大小。我习惯是 2GB 起步大数据量场景给到 4GB但不要超过物理内存的一半毕竟 ES 本身也是吃内存大户。还有几个 pipeline 参数在config/logstash.yml里pipeline.workers: 4 pipeline.batch.size: 500 pipeline.batch.delay: 50pipeline.workers是执行 filter 和 output 的线程数建议不超过 CPU 核数设太高反而增加上下文切换开销。pipeline.batch.size控制每个 worker 一次性批量处理的事件数默认 125可以调到 500 或 1000。但别贪大批量越大单次 bulk 请求也越大ES 处理不过来时容易超时重试反而拖慢整体速度。我调参的原则是小步试观察日志里 bulk 请求的耗时和失败率再决定要不要继续加。6.2 控制查询压力SQL 层面的优化很多人只盯着 Logstash 调优忽略了数据库才是瓶颈。JDBC 插件每 5 秒执行一次全表扫描MySQL 迟早被拖垮。SQL 层面有几个必须做的优化增量字段必须建索引。WHERE update_time :sql_last_value这种查询如果update_time上没有索引每次都是全表扫描数据量一大数据库 IO 直接打满。这是同步场景最常见的问题。避免SELECT *只查需要的列。尤其是 TEXT/BLOB 类型的大字段拖慢查询速度还涨 ES 索引占用。增量查询加时间上界。建议写成SELECT id, order_no, user_id, status, update_time FROM orders WHERE update_time :sql_last_value AND update_time NOW() - INTERVAL 1 MINUTE加这个上界的意义在于正在被业务更新的数据可能刚改一半同步过去的是中间状态。加个 1 分钟的时间窗口让数据稳定了再同步可以避免频繁覆盖写入。缺点是多了一分钟延迟对大多数搜索场景来说完全能接受。大表全量不要一条 SQL
上一篇/下一篇内容由系统自动关联 返回资讯列表 →