尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

MySQL实时同步Elasticsearch:Logstash全量与增量实践

MySQL实时同步Elasticsearch:Logstash全量与增量实践 先给结论只要不是那种每秒几万次写入的极端场景MySQL到Elasticsearch的实时同步最稳、最省事的方案就是让Logstash干这个活。我在云主机上折腾过好几套方案最后留在生产环境的也是这一套。整个过程踩了不少坑从全量导入到增量监听从字段类型映射到删除数据同步这篇文章把能想到的问题都梳理一遍希望能帮你少走弯路。同步这件事看起来简单不就是把数据库的数据搬过去嘛但真做起来会发现到处是细节全量同步怎么不把服务压垮、增量同步怎么判断哪些数据改了、删除操作怎么让ES也删掉、数据结构变了怎么办。这篇文章会从架构设计讲到具体配置再附上运行中经常会遇到的坑完整走一遍MySQL实时同步Elasticsearch的流程。1. 为什么需要同步以及同步方案怎么选1.1 业务里最常见的两个诉求我接到这类需求时一般先问两个问题一是要拿这些数据做什么二是对“实时”能接受多少延迟。这两个问题直接决定方案走向。如果是把MySQL里的商品、订单、文章等内容同步到ES做搜索那通常只需要秒级延迟Logstash轮询就够用。如果是金融风控、实时大屏这类场景延迟要控制在毫秒级那就得用基于日志捕获的方式来监听数据变更比如Canal、Flink CDC这些。架构上差别很大成本也不一样。Elasticsearch擅长的是全文检索、复杂聚合、大数据量下的快速查询MySQL则擅长事务性写入和关系型查询。两者天然互补。但这里的“互补”是建立在数据一致性上的如果ES里的数据和MySQL对不上搜索出来的结果就是错的所以同步链路的可靠性比同步本身更重要。1.2 方案选型不是越高级越好把市面上常用的几种方案放在一起对比过之后你会发现没有标准答案只有适不适合方案同步原理适合场景维护成本延迟Logstash JDBC插件定时轮询SQL查询中小规模、延迟容忍秒级低秒级Canal 客户端订阅MySQL binlog需要监听删除、低频高实时中高毫秒级Flink CDC全量增量一体化大数据量、复杂计算高毫秒级DataX批量离线导入一次性全量迁移低分钟级一开始我也被“必须毫秒级”的念头带跑过后来冷静下来算了一笔账数据量几千条级别的同步Logstash轮询间隔配3秒成本最低、配置最简单出问题也好排查。Flink CDC确实高级但要部署Flink集群要写SQL要处理checkpoint对于一个小团队来说运维负担很重。Canal这个方案我在中间件场景用过它监听binlog很精准删除操作也能捕获到不会出现ES里残留脏数据的问题。但代价是Canal Server本身要部署客户端代码要自己维护还要处理binlog解析、位点续传。如果你的团队没有专门的中间件开发人力我建议还是先用Logstash把核心链路跑通等真遇到性能瓶颈再升级。1.3 我的最终选择Logstash为主全量和增量分开做最终我敲定的方案是全量同步用Logstash跑一次大查询增量同步用Logstash定期轮询增量字段。为什么不全量增量用一个管道因为两者的SQL、字段处理策略、对数据库的压力完全不一样混在一起容易把问题搞复杂。全量阶段我会用一个单独的管道文件一次性把表里的所有数据拉出来写入ES。这个阶段重点是控制查询节奏避免一个大SQL把数据库拖垮。增量阶段我会把同步时间窗口记录下来只查上次同步之后发生变更的数据。这个阶段重点是准确识别变更包括新增和修改。这套组合的好处是部署简单只需要安装一个Logstash逻辑清晰全量和增量互不干扰排查方便日志里每一步都能看到同步了多少条数据、用了多少时间。对于大部分中小型项目来说这已经是足够优秀的方案了。2. 实时同步的核心机制从全量快照到增量监听2.1 先搞清楚MySQL到ES的数据流如果你第一次接触同步脑子里先得有这条链路MySQL里的一张表经过查询、转换、填充最终变成ES里的一条文档。这个“文档”的ID通常就是MySQL的主键文档内容就是这一行的字段组合。不知道你有没有用过ES的_id字段它在文档写入时可以指定Logstash也能把MySQL的主键直接映射成文档ID。这么做的好处是幂等同一行数据重复同步多少次ES里都只有一条记录。如果不指定IDES会为每条文档生成随机ID重复同步就会产生大量重复数据后患无穷。增量同步之所以能实现核心就靠两个东西一个是MySQL表里的时间字段比如updated_at另一个是Logstash执行查询的时间记录也就是sql_last_value。每次同步时Logstash会跑一条带条件的大于上次时间戳的SQL把新产生或修改过的数据取出来写进ES然后再把当前时间记录下来作为下一次查询的起点。就这么简单。2.2 binlog能不能直接用来同步这里要聊一下更底层的binlog机制。MySQL的binlog记录了数据库的所有变更操作包括INSERT、UPDATE、DELETE。Canal和Flink CDC之所以能做到低延迟就是因为它们能直接读取并解析binlog拿到每一条变更的原始数据。但binlog不是你想读就能读的首先要开启MySQL的binlog配置其次要确保binlog格式是ROW因为只有行级格式才会记录每一行数据的变更内容最后你还要处理binlog的过期清理和位点保存不然服务重启后不知道从哪断点继续。这套东西玩好了确实牛但玩不好就是灾难。我在生产环境见过一次binlog被清理后Canal同步直接从崩掉的位点接着读结果全量数据重新导了一遍幸好当时数据量不大。如果数据量到了几十G的量级这种事情会让人非常崩溃。所以我的原则是能用字段轮询解决的就别太早碰binlog。2.3 删除数据怎么处理这是个灵魂问题用Logstash的JDBC插件做增量同步时只能查到“还存在”的数据MySQL里被DELETE掉的记录ES并不会有任何感知。也就是说你用这种方式同步的ES索引会一直留着那些已经删除的文档只有重建索引才能清掉。这个问题怎么解决几个常见方法第一如果业务允许不要物理删除改用逻辑删除加一个deleted字段。同步时只查deleted0的数据ES里自然不会出现被删掉的记录。第二定期做全量重建比如每天凌晨跑一次全量同步覆盖当天产生的小部分脏数据保证索引长期可用。第三接入Canal或Flink CDC从binlog里拿到DELETE操作主动把ES里对应的文档删掉。要彻底解决删除同步问题最后绕不开这条路。我在实际项目里通常的做法是先用逻辑删除顶着同时把全量重建做成一个定时任务。等业务量真的上来了再考虑引入Canal。这个演进节奏很关键别一开始就上重型方案。3. 实操用Logstash完成全量与增量同步3.1 环境准备与参数规划在动手配置前先列一下我的环境参考值云主机4核8G跑Logstash和Elasticsearch同一台机上MySQL版本5.7和8.0都试过配置方法一致Elasticsearch版本7.x和8.x都兼容推荐8.xLogstash版本我用的是7.17.xJDBC插件自带JDBC插件的安装这里不重复说了需要注意一点Logstash默认自带logstash-input-jdbc这个插件如果是老版本需要手动安装命令是bin/logstash-plugin install logstash-input-jdbc。JDBC驱动别用老的5.1版本MySQL 8.0必须用mysql-connector-java-8.0.x.jar放进Logstash的某个目录然后在配置里用jdbc_driver_library参数指定路径。很多同步异常都是驱动版本和数据库版本对不上导致的尤其是SSL连接报错、字符集乱码基本都是驱动问题。3.2 全量同步的配置模板全量同步的核心是把整张表的数据捞出来写进ES。配置看起来长其实核心就是input、filter、output三件事。先给一个能直接改的表结构做示例。假设有一张t_article表CREATE TABLE t_article ( id bigint NOT NULL AUTO_INCREMENT, title varchar(200) DEFAULT NULL, content text, author varchar(50) DEFAULT NULL, status tinyint DEFAULT NULL, created_at datetime DEFAULT NULL, updated_at datetime DEFAULT NULL, PRIMARY KEY (id) ) ENGINEInnoDB;对应的Logstash管道配置存成全量同步专用的配置文件比如sync_full.confinput { jdbc { jdbc_driver_library /usr/local/logstash/lib/mysql-connector-java-8.0.26.jar jdbc_driver_class com.mysql.jdbc.Driver jdbc_connection_string jdbc:mysql://127.0.0.1:3306/test_db?useSSLfalsecharacterEncodingutf8 jdbc_user sync_user jdbc_password your_password jdbc_paging_enabled true jdbc_page_size 5000 statement_filepath /usr/local/logstash/config/full_query.sql schedule * * 3 * * } } filter { mutate { remove_field [version, timestamp] } } output { elasticsearch { hosts [http://127.0.0.1:9200] index article_index document_id %{id} pipeline article_pipeline } }需要注意上面的schedule我用的是每天凌晨3点执行这样全量同步就不用人工触发直接由定时任务驱动。平时在测试时也可以单独执行一次把schedule参数去掉或注释掉然后启动Logstash跑完即退。这里最关键的是document_id必须设置。如果不设置每次全量同步都会把同一批数据当成全新文档写入ES造成相同ID的数据存在多份写入性能差而且数据量大时磁盘直接爆掉。设置成%{id}以后每次同步相当于UPSERT存在就更新不存在就新增非常稳。statement_filepath指向一个独立的SQL文件这样SQL和配置分离以后改动查询逻辑不用动整个管道配置。我的全量SQL通常这样写SELECT id, title, content, author, status, created_at, updated_at FROM t_article全量同步不需要WHERE条件但如果你只想同步某一部分数据也可以在这里加WHERE过滤。写入ES时Logstash会自动把MySQL的字段类型转成JSON类型比如datetime会变成字符串bigint会变成数字。这些在ES端都要提前把mapping建好否则ES会按默认规则猜测字段类型后面就麻烦了。3.3 增量同步的配置模板增量同步和全量同步的核心区别在SQL。如果把全量同步比作“把仓库里的货全盘一遍”增量同步就是“只把今天新进的货记下来”。增量配置如下input { jdbc { jdbc_driver_library /usr/local/logstash/lib/mysql-connector-java-8.0.26.jar jdbc_driver_class com.mysql.jdbc.Driver jdbc_connection_string jdbc:mysql://127.0.0.1:3306/test_db?useSSLfalsecharacterEncodingutf8 jdbc_user sync_user jdbc_password your_password jdbc_paging_enabled true jdbc_page_size 1000 statement_filepath /usr/local/logstash/config/increment_query.sql schedule */3 * * * * * use_column_value true tracking_column unix_ts tracking_column_type numeric record_last_run true last_run_metadata_path /usr/local/logstash/data/.logstash_jdbc_last_run } } output { elasticsearch { hosts [http://127.0.0.1:9200] index article_index document_id %{id} pipeline article_pipeline } }这里面的货比较多我给你一个个说明白tracking_column是Logstash用来比较数据的参考字段。它不一定是时间字段只要是单调递增就行。如果你有自增ID也可以用ID做追踪。不过ID做追踪的话就检测不到UPDATE操作因为UPDATE之后ID不变。所以最稳妥的做法是用时间字段。tracking_column_type我用的numeric配合unix_ts这样一个数字类型的字段使用。为什么不直接用datetime类型做追踪因为MySQL返回的datetime格式有时候带有时区差异字符串比较还会出问题。把所有时间统一转成Unix时间戳再比较是最省心的做法。在写增量查询SQL时需要在SQL里动态取出上次同步的位置用:sql_last_value这个占位符。我的increment_query.sql长这样SELECT id, title, content, author, status, created_at, updated_at, UNIX_TIMESTAMP(updated_at) AS unix_ts FROM t_article WHERE UNIX_TIMESTAMP(updated_at) :sql_last_value ORDER BY updated_at ASC LIMIT 2000关于LIMIT这里要特别提一下。如果你不在SQL里做LIMITLogstash会在查询到所有结果之后做内存分页大批量数据时内存占用非常高。我在数据量超过10万条时明显感觉到Logstash内存飙升后来在SQL里加上合适的LIMIT再配合jdbc_page_size参数内存占用就稳定下来了。但这里有个坑要注意如果加LIMIT必须加ORDER BY updated_at ASC否则系统可能把“最早更新”的数据留在后面而同步进度却已经往前跳了导致漏数据。ORDER BY LIMIT组合起来才能保证每次都从最早更新的数据开始处理这是我自己踩过坑才注意到的。3.4 同步过程中的函数与过滤逻辑从MySQL读出的数据未必直接就能写入ES。常见情况有三种第一种是字段类型不匹配。比如MySQL的十进制字段传到ES后可能变成字符串再比如JSON字段MySQL拿到的可能是个字符串需要用json过滤器解析成JSON对象ES才能参与嵌套查询。第二种是字段名冲突。MySQL字段名里的下划线和驼峰风格问题经常让人头疼。我习惯在filter里统一改成ES的命名习惯比如把author_name改成authorName。第三种是多余字段需要剔除。JDBC插件默认会把version、timestamp等Logstash内部字段加进去如果不想让它们在ES文档里出现就在filter里用mutate把默认字段删掉。我自己在filter里会写这样一段filter { mutate { rename { author authorName created_at createTime updated_at updateTime } remove_field [version, timestamp] } json { source content target contentObj remove_field [content] } }这段逻辑是把content的JSON字符串解析成contentObj这样ES里就能对文章内容做结构化了。如果你不需要结构化直接保留字符串文本也可以。3.5 给Logstash配置内存和运行参数Logstash默认分配1G堆内存同步大数据量时根本不够。这个我深有体会。全量同步跑到一半Logstash直接OOM崩溃看起来像是ES写入失败其实根因是JVM堆内存爆了。我推荐的参数设置是# 修改 Logstash/config/jvm.options -Xms2g -Xmx2g如果你的机器有8G内存给Logstash分2G是合理的再多的话会挤压ES的JVM内存。ES的JVM内存建议不要超过物理内存的50%因为ES还要用额外的堆外内存做文件缓存。两头都不能贪婪要平衡。运行Logstash时我会用-r参数指定配置目录cd /usr/local/logstash bin/logstash -f config/sync_full.conf如果想让Logstash在后台运行可以用nohup或者配置成systemd服务。生产环境我建议用systemd管理这样进程挂了能自动拉起。systemd里要设置Restartalways和RestartSec5避免因为一个网络抖动就导致同步永久停止。4. 生产环境经验常见问题与避坑速查4.1 同步延迟怎么排查同步延迟是实时同步的核心指标。我一般通过curl查ES里的最新文档时间来判断延迟curl -X GET http://127.0.0.1:9200/article_index/_search -H Content-Type: application/json -d { size: 1, sort: [{ updateTime: desc }] }看返回的updateTime和当前时间的差值如果超过1分钟说明Logstash的轮询频率不够或者查询效率有问题。先看日志Logstash启动后会输出每次JDBC查询的结果条数和耗时。一般来说延迟大是下面几个原因SQL没有走索引全表扫描耗时太长。给updated_at字段加上索引是最常见的优化手段。每次查询的数据量太大Logstash处理不过来。调整LIMIT大小和jdbc_page_size。ES写入端出现节点压力bulk队列堆积。关注ES的写入线程池指标。4.2 MySQL时区导致的8小时时间差刚部署完同步时我遇到过一个非常诡异的问题ES里的时间和MySQL里的时间差了8个小时。排查了一圈确认是时区配置不一致。Logstash的JDBC连接串里没有指定时区时默认使用JVM所在地时区而MySQL的连接时区又是另一个值。解决方法是统一时区配置。在JDBC连接串上显式加上jdbc:mysql://127.0.0.1:3306/test_db?useSSLfalsecharacterEncodingutf8serverTimezoneAsia/Shanghai同时在Elasticsearch的索引mapping里日期字段建议存储为yyyy-MM-dd HH:mm:ss字符串格式或epoch_millis数字格式不要在代码层面搞时区转换。最省事的是直接用UNIX_TIMESTAMP转成数字存进ES这样完全不受时区影响。4.3 字段类型冲突与ES Mapping规划ES是schema-free的但这不代表你能随意写入。第一次写入某字段时ES会猜测字段类型后面如果类型对不上就会报错。比如MySQL的bigint转成ES的long没问题但如果int类型字段在mapping里被定义成了text写入整数时ES会报mapper_parsing_exception。开始全量同步前先把索引mapping建好这个步骤不能省。我的做法是提前用POST请求创建索引curl -X PUT http://127.0.0.1:9200/article_index -H Content-Type: application/json -d { mappings: { properties: { id: { type: long }, title: { type: text, analyzer: ik_max_word }, contentObj: { type: object }, authorName: { type: keyword }, status: { type: integer }, createTime: { type: date, format: yyyy-MM-dd HH:mm:ss||epoch_millis }, updateTime: { type: date, format: yyyy-MM-dd HH:mm:ss||epoch_millis } } } }title字段如果要支持中文分词需要安装IK中文分词插件。这里一个小提醒mapping一旦创建普通字段类型就不能热修改了。如果字段类型建错了最干脆的办法是删除索引重建别想着ES能给你改类型它不能。4.4 数据库压力过大怎么办全量同步时如果一次性SELECT *一个大表数据库很容易被打满尤其在生产环境。我的经验是把查询分段跑利用MySQL的WHERE id ? AND id ?做分段控制或者靠LIMIT分页。Logstash的jdbc_paging_enabled参数就是干这个的开启后它会自动分页查询避免一次加载全部结果到内存。另一个更稳的做法是全量同步放在业务低峰期执行比如凌晨。如果业务全天都不能接受明显压力就要考虑用主从库的从库作为同步数据源从库只读对主库毫无压力。这是大项目里普遍采用的方案。4.5 数据重复和漏数据怎么处理重复数据的根源通常是重复执行同步任务比如多个Logstash进程同时跑同一个管道。检查方法很简单ps -ef | grep logstash看看进程数。Logstash的管道自己会起多个worker线程但同一个配置不要用多个进程跑。漏数据则多见于增量SQL的查询边界问题。如果你用where updated_at :sql_last_value且恰好有一行数据的updated_at时间与上次记录的时间完全相同这行数据就可能被漏掉。解决办法是用再配合主键去重或者让更新时间精确到毫秒。时间字段精度越高这种边界重合的概率越低。我在生产环境遇到过一个小概率事件某一行数据的updated_at被回滚成过去的时间导致这行数据永远不会进入增量同步范围。如果业务上确实存在这种操作那就只能安排每日全量重建来兜底。4.6 ES和Logstash版本兼容问题Logstash和Elasticsearch版本跨度不能太大官方约定是大版本必须一致。比如ES是7.xLogstash就用7.xES是8.xLogstash就用8.x。我见过有人ES升到8.xLogstash还在6.x结果同步时出现各种诡异的序列化错误。从ES 8.0开始默认开启了安全认证Logstash连接ES时需要在output里配置用户名密码output { elasticsearch { hosts [https://127.0.0.1:9200] user logstash_system password your_password ssl true } }同时ES的xpack.security.enabled要保持启用状态。很多人在开发环境图省事关掉安全认证结果部署到生产环境忘了开连接直接失败。最后再分享两个小经验这套同步链路跑起来之后我最大的体会是别迷信“一键同步”这种黑科技实时同步的本质就是处理数据变化。你能在多短的时间内发现数据变了、多可靠地把这个变化搬到目标端决定了整个链路的稳定性。实际运营中我每天都会盯一眼同步延迟每周核对一次ES文档数和MySQL表行数。如果数量对不上不一定是同步逻辑有问题可能是MySQL有软删除、ES有删除策略但不管怎样及时发现总比用户反馈要强。还有一个扩展思路当Logstash扛不住数据量时不要急着换Flink先试试把同步任务拆散成多个管道按库分、按表分、按字段分并行度提上来很多性能问题就解决了。真正到了单管道、多管道都撑不住的时候再考虑引入Canal或Flink CDC不迟。这样每一步都有明确的回报踩的坑也会少很多。
返回列表