Canal实战:基于Binlog实现MySQL到Elasticsearch增量数据同步
先聊个实在的很多团队一开始做MySQL到ES的数据同步第一反应就是自己写定时任务凌晨跑一次全量或者用Logstash做增量轮询。这两种方式要么有延迟要么实现复杂数据量大之后基本都会出问题。后来我接触了canal才意识到“监听MySQL binlog、实时推给ES”才是增量同步该有的样子。这篇文章就把我实际部署canal同步MySQL增量数据到ES的完整过程、方案取舍和踩过的坑一次讲清楚适合正在做搜索引擎、BI报表、复杂查询加速的开发和运维同学参考。canal是什么简单说它是阿里巴巴开源的一个中间件能够把自己伪装成一个MySQL从库通过MySQL主从复制协议去订阅binlog然后把binlog里的增删改事件解析出来再交给下游消费。配合canal adapter我们可以把解析结果直接写入ES、Kafka、RocketMQ等目标端实现近乎实时的数据同步。在ES场景里这套链路解决的是“业务库负责写ES负责读”的典型异步架构问题也是目前增量同步到ES的主流方案之一。不管你是刚接触canal的新手还是已经踩过一些坑的实践者这篇文章都值得先收藏再往下看。下面我会从方案设计、环境准备、配置实现到故障排查按实际落地的顺序一步步拆开。1. canal同步方案的设计思路与核心原理1.1 为什么业务侧需要“MySQL到ES”的增量同步先明确一个背景MySQL擅长事务和强一致性但遇到模糊搜索、多条件组合筛选、统计聚合这类查询性能就会快速劣化。ES则天生擅长倒排索引和分布式搜索。所以很多系统会同时保留两份数据MySQL管写ES管查。这带来一个新问题MySQL里的数据变了ES里的数据怎么跟着变常见做法是业务代码双写发一条消息到MQ消费端再写ES。听着简单但侵入性很强每个业务模块都得改代码。最重要的是一旦某条消息丢了或者顺序错乱MySQL和ES的数据就永久不一致问题定位起来极其痛苦。另一种做法是定时全量重灌数据量小的时候还能糊弄过去到了百万、千万级别跑一次全量要十几分钟甚至更久查询期间数据是旧的就忍忍吧更别提每天晚上都得占机器资源。对比之下canal做的增量同步是从数据库日志层面获取变化业务代码完全不用改不侵入业务链路。只要MySQL开了binlogcanal就能实时拿到所有数据变更想写ES就写ES想入Kafka就入Kafka。这是它成为主流的根本原因。1.2 canal的工作机制把自己伪装成MySQL从库这里有必要把原理讲透。MySQL主从复制本身是这么运作的主库开启binlog后所有数据变更都会以二进制日志形式记录下来从库连接主库把自己伪装成一个服务器ID不同的普通客户端然后向主库请求binlog主库按位置推送日志从库拿到后本地回放从而实现数据复制。canal做的事情和“从库请求binlog”几乎一样。客户端向canal发起订阅请求canal服务端从MySQL拉取binlog解析出结构化的数据变更事件再以自有的client协议推给下游。官方自带的canal.adapter是一个开箱即用的对接组件支持把变更事件直接写入ES、HBase、Kafka等目标端。我们下面讨论的链路就是MySQL binlog - canal server - canal adapter - Elasticsearch。需要注意的是canal只是搬运工它不负责ES索引的full text搜索逻辑。ES这边的索引结构、分词器、字段映射仍然需要你在同步前先设计好。整个同步方案里一半是canal的配置另一半是ES索引和映射的设计两者缺一不可。1.3 方案选型adapter直连ES还是走MQ再消费在确认要用canal之后很多人会纠结是直接用adapter直连ES还是在canal和ES中间再插一层MQKafka/RocketMQ。我的建议很直接同步场景简单、并发量可控、团队规模不大直接走adapter省事维护成本低。但如果你的下游不止ES一个后面还要接数仓、缓存、其他搜索集群或者ES写入压力需要削峰填谷那就必须引入MQ。直连ES的优势是链路短、延迟低canal server解析出增量事件后adapter直接打包成bulk请求写入ES配置一次就能跑。缺点是当ES写入抖动或者集群过载时adapter会因为写ES太慢而积压进而产生数据同步延迟。接入MQ的做法是把消费端和canal解耦adapter只负责把变更事件投递到Kafka指定topic后端的消费程序根据自己的节奏消费并写入ES。这样即使ES抖动也只是MQ消费停滞canal仍然可以继续从MySQL拉日志不会阻塞主库。缺点是需要多维护一套MQ集群和消费程序配置链路也长一截。我实际项目里比较推荐的组合是初期用adapter直连ES快速跑通同步稳定后如果发现ES写入流量太大或者需要多下游消费再平滑迁移到“canal Kafka 自研消费端”。这样不会一上来就被复杂的组件拖住。2. 部署落地前的环境准备2.1 MySQL端必然要准备的三个基本条件canal要读取MySQL的binlog第一步就是确保MySQL开启了binlog官方推荐的配置格式是row。为什么必须是row因为statement格式记录的是SQL语句本身在同步端执行同样语句可能因为环境差异产生不同结果row格式记录的是每一行数据“长什么样”改动前和改动后都有canal才能准确地把整行数据同步到ES。这里直接给一份最小可用的MySQL配置[mysqld] log-binmysql-bin binlog-formatROW server-id1 # 如果用了GTID下面这个打开更利于后续管理 gtid_modeon enforce_gtid_consistencyon配置修改后记得重启MySQL。然后创建canal专用的同步账号不需要给超级权限但要给复制权限CREATE USER canal% IDENTIFIED BY canal_pass; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;这里有个容易踩的坑如果MySQL的server_id设置了1而canal连接时也默认给server-id 1就会和主库或者其他从库冲突导致canal连接失败或者复制中断。建议canal的server-id单独设置一个不冲突的值比如1到2的32次方之间的随机数。2.2 canal server部署两种方式按团队情况选择canal server是整条链路的枢纽负责解析binlog和分析事件。部署方式我试过两种直接下载压缩包启动以及用Docker跑容器。压缩包方式排查问题方便日志文件路径清楚Docker方式适合快速验证和环境隔离。这里先讲压缩包方式因为很多生产环境不让随便跑容器。去canal官方GitHub Releases页面下载对应版本的canal.deployer压缩包比如canal.deployer-1.1.7.tar.gz。解压后核心配置在conf/canal.properties里canal.port11111 canal.destinationsexample canal.instance.global.modemanagercanal server启动后默认监听11111端口下游adapter通过这个端口连接。每个instance实例对应一个MySQL主库连接默认的example实例在conf/example/instance.properties里配置具体的MySQL连接信息canal.instance.master.address127.0.0.1:3306 canal.instance.dbUsernamecanal canal.instance.dbPasswordcanal_pass canal.instance.connectionCharsetUTF-8 canal.instance.tsdb.enabletrue启动canal servercd canal.deployer-1.1.7 sh bin/startup.sh日志输出在logs/example/example.log如果看到类似“dump address ... has been successfully found”的提示说明已经成功连接上MySQL了。Docker方式更简单一条命令就能起一个serverdocker run -d --name canal-server -p 11111:11111 canal/canal-server:latest但容器里的配置一般通过环境变量或者挂载配置文件来做更适合你已经熟悉canal配置之后再操作。第一次上手建议还是压缩包日志看起来更方便。2.3 canal adapter部署与基础配置adapter是真正对接ES的那一层。同样去官方下载canal.adapter压缩包解压后你会看到conf目录下有两份核心配置application.yml和bootstrap.yml通常只需要改application.yml。先把adapter和server的连接串起来canal.adapter.manager: - canal: server: 127.0.0.1:11111 destination: example username: password:然后是ES连接配置。adapter默认使用RestHighLevelClient风格的接口连接ES需要指定ES集群地址、认证信息。比较关键的几个参数包括es: hosts: http://127.0.0.1:9200 username: elastic password: elastic_pass这里有个容易忽略的细节adapter连接ES的方式在1.1.5之后有所调整如果ES开启了安全认证必须确认配置文件里的用户名密码正确否则启动时不会报错但同步写数据时会疯狂刷认证失败的异常效果特别迷惑。bootstrap.yml里面主要配置的是spring应用相关项比如gin模式、数据源连接等一般不用动。只有当你的adapter不止一台并且用了ZK或者Manager模式做集群协调时才需要去改这些。启动adaptersh bin/startup.sh正常启动后adapter会去连接canal server如果你同时打开canal server的日志和adapter的日志会发现两边建立订阅关系的全过程。这一步通了后面的同步配置才有意义。3. 同步任务配置与核心实现细节3.1 同步映射文件决定MySQL数据如何落到ES索引canal adapter对“哪张表同步到哪个索引”这件事完全是靠conf/es/目录下的一批.yml映射文件控制的。每个映射文件定义了一个同步任务比如orders.yml就表示把MySQL里的orders表同步到ES的orders索引。一个最基础的映射文件长这样dataSourceKey: defaultDS destination: example groupId: g1 esMapping: _index: orders _id: id sql: SELECT id, order_no, amount, user_id, status, create_time, update_time FROM orders commitBatch: 3000这里有几个关键点需要说清楚。dataSourceKey对应application.yml里配置的数据源名称默认是defaultDS。destination要和你canal server里的instance名称保持一致一般就是example。_index是ES索引名注意这里不能用大写的索引名ES索引名规范要求全小写。_id怎么选很关键最好和MySQL表主键或业务唯一键保持一致这样每次更新都是基于同一文档ID做upsert不会产生重复文档。sql字段是adapter做数据转换的核心它决定了下游ES里一条文档对应MySQL的哪些字段。注意一个坑adapter执行这个sql时会把canal解析出来的主键值作为查询条件所以sql里最好显式包含where条件依赖的主键字段。默认adapter会自动拼上主键id条件但如果你在sql里对字段做了别名或者别名太复杂会导致自动拼接失败。我习惯写sql就是select原表字段保持简单。另外如果你的ES索引里有一些字段在MySQL里不存在比如一个固定的type字段或者一个用于过滤的category_code你可以用sql别名的方式生成。adapter并不会限制你只能select表的原始列它拿到的sql运行结果就是最终写入ES文档的内容。3.2 增量同步验证从insert到update再到delete配置完成后激动人心的时刻就是验证增量同步。先在MySQL里执行一条insertINSERT INTO orders (id, order_no, amount, user_id, status) VALUES (1001, NO20240111001, 199.00, 88, PAID);然后立即查询ESGET orders/_doc/1001正常情况下几秒内就能看到这条文档出现在ES里。整个过程中canal解析binlog、adapter执行sql、再bulk写入ES链路里的每一步都有日志可查哪一步慢了或者断了都能在日志里找到痕迹。接下来测updateUPDATE orders SET status SHIPPED, update_time NOW() WHERE id 1001;再去ES查询会发现文档里的status字段已经被更新。这就是增量同步的核心价值业务侧完全不需要感知ES的存在数据却实时保持一致。最后测deleteDELETE FROM orders WHERE id 1001;这里要特别提醒adapter默认的删除映射逻辑是根据主键_id去ES里删除文档如果你的MySQL表有物理删除同步到ES就是delete by id。但很多系统的MySQL并不做物理删除而是用status字段标记删除比如status0表示已删除。这种情况你在ES里定义文档时推荐同步status字段并在查询端加过滤器而不是让canal去删除doc因为canal只在row模式下能拿到物理删除事件逻辑删除对它来说就是一次普通update。3.3 存量数据怎么同步先全量再增量的衔接canal的核心是增量但它不解决存量数据初始化问题。第一次做同步的时候ES索引里可能一条数据都没有这时如果只开增量历史数据永远不会进ES。正确的做法是先做一次全量导入再开启增量同步两者衔接好避免中间漏数据。全量导入我常用的工具是canal adapter自带的全量同步功能。构造一个http请求curl -X POST http://127.0.0.1:8081/sync?destinationexamplegroupg1这个接口会读取映射文件中的sql执行一次全量select然后把所有结果写入ES。执行完之后再确保canal server是从当前binlog位点开始监听增量。注意事项是全量同步的耗时和MySQL表数据量直接相关大表全量期间产生的增量变更理论上会写入binlog并在之后被canal消费到所以整体不会丢但要观察ES里的数据是否最终一致。更稳妥的做法是在业务低峰期做全量结束后人工抽样对比MySQL和ES的count再开启增量。如果中途出现不一致不要急着手动往ES补数据先查binlog位点和adapter日志定位是同步链路断了还是同步配置有误。3.4 字段类型与时间时区同步最容易翻车的地方ES和MySQL的数据类型不是一一对应的。MySQL里的decimal、datetime、json在ES里都要提前规划和映射。最简单的方式是让ES自动映射查String字段时用keyword类型查数值时用long或double查时间自动识别为date。但生产环境我不推荐全自动映射因为你无法控制分词行为连es的keyword和text都分不清后面查询一定会后悔。以我的经验同步映射文件里至少要显式处理两种类型。一种是时间字段MySQL的datetime和timestamp通过adapter同步到ES后通常被识别为date类型。ES默认date格式是ISO8601类似2024-01-11T12:00:00Z。如果你的MySQL时间字段是字符串或者期望用时间戳就要在ES映射里明确format{ mappings: { properties: { create_time: { type: date, format: yyyy-MM-dd HH:mm:ss||yyyy-MM-ddTHH:mm:ss.SSSZ||epoch_millis } } } }另一种是MySQL的json类型。json字段在MySQL里比较灵活但ES里不能直接映射为object自动搞定尤其是json里嵌套了数组或者动态key的建议在ES里指定为flattened类型或者干脆把json序列化成字符串存到text字段由业务查询时自行解析。别偷懒把所有json都映射成nested索引变大且维护成本高。时区问题更要命。canal在解析binlog时time类型输出的是UTC时间或本地时间取决于instance.properties里的canal.instance.connectionCharset和默认时区配置。如果你在MySQL里存储的CST时区的时间但canal解析后转换成了UTC同步到ES里会差8小时。排查这个问题的思路是先看MySQL原始时间再看canal解析日志里的时间最后看ES文档里的时间对比出是在哪一步出现了偏移再去调整时区配置不要一上来就乱改ES映射。4. 实战操作步骤从零搭一套MySQL到ES的增量同步4.1 步骤一确认MySQL binlog已经开启并且格式为row进入MySQL命令行执行SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format;如果log_bin是OFF或者binlog_format不是ROW就得去my.cnf里改配置并重启MySQL。这一步是绝对前提canal没binlog可读后面全部白搭。另外建议顺便执行SHOW MASTER STATUS;记录当前binlog文件名和位置这个信息在后续验证增量位点时很有用。如果是刚开启binlog的实例第一次需要让canal从当前位点开始监听不要从老位点开始否则可能找不到了。4.2 步骤二部署canal server到可用状态假设MySQL跑在192.168.1.100:3306canal server部署在同网段的另一台机器。按前文的配置修改canal.properties和instance.properties然后启动。启动后看到一个关键日志“start successful”就说明main进程起来了但还要确认instance是否连接了MySQL。日志里出现类似“begin find start position”的时候说明canal已经成功从MySQL拿到了binglog消费位点。这里给你的一个良好习惯把canal server的logs目录单独用软链指到有足够硬盘空间的路径因为canal在某些模式下降了批量的解析日志时间长了会占不少空间。别等磁盘满了才想起来清理很被动。4.3 步骤三部署adapter并连接canal server解压adapter后修改conf/application.yml把canal server地址和ES地址都填好。如果MySQL连接信息不止一个数据源可以在这个配置文件里定义多个不同key的数据源映射文件用dataSourceKey区分。我觉得这个设计其实很实用比如你有多套库表要往同一个ES集群同步不用部署多套adapter一个adapter就能配置多个同步任务。启动后观察adapter日志正常时会出现“canal adapter started”和“subscribe ... success”之类的信息。一旦出现订阅失败先把canal server的日志打开看看是不是server端本身就和MySQL断开了再检查网络和防火墙。4.4 步骤四编写ES索引映射并创建索引canal adapter不会自动创建ES索引虽然有些版本在同步时会用默认映射创建索引但字段类型往往不是你想要的。所以我们要在启动同步之前先把索引建好。这里以orders索引为例PUT orders { settings: { number_of_shards: 3, number_of_replicas: 1 }, mappings: { properties: { id: { type: long }, order_no: { type: keyword }, amount: { type: scaled_float, scaling_factor: 100 }, user_id: { type: long }, status: { type: keyword }, create_time: { type: date, format: yyyy-MM-dd HH:mm:ss||yyyy-MM-ddTHH:mm:ss.SSSZ||epoch_millis }, update_time: { type: date, format: yyyy-MM-dd HH:mm:ss||yyyy-MM-ddTHH:mm:ss.SSSZ||epoch_millis } } } }注意amount用scaled_float这样小数精度能保留且ES内部是long存储查询统计更高效。time字段在ES里必须要有统一的格式否则不同的写入值会parse异常bulk请求直接失败。索引创建好之后再启动adapter或后写映射文件数据进来就能顺利落盘。4.5 步骤五做一次完整的增删改验证索引建好adapter也跑起来了接下来做完整验证。先在MySQL里写入多条测试数据覆盖各种类型字符串、整数、小数、日期。然后等几秒查询ES确认写入。再把其中一条记录update成新值确认ES里的文档内容跟着变了。最后delete一条确认adapter的删除事件生效ES这条文档也没了。一个容易被忽视的验证点如果你在MySQL里的update没有修改任何实际内容比如set的值和原值一样MySQL会不会产生binlog不会。那canal自然不会感知到ES也不会产生更新。这不算故障是合理行为不用排查方向搞偏了。4.6 步骤六启动守护进程或systemd服务canal server和adapter都是Java进程生产环境不可能开着终端跑。我一般用systemd来守护这样崩溃了能自动拉起。给canal server写一个简单的service文件[Unit] Descriptioncanal server Afternetwork.target mysql.service [Service] Usercanal ExecStart/opt/canal/bin/startup.sh ExecStop/opt/canal/bin/stop.sh Restartalways RestartSec5 [Install] WantedBymulti-user.targetadapter同理。有一点要注意canal的startup.sh脚本内部会通过jvm参数启动如果你用systemd强杀进程可能留下状态文件让下次启动异常。我踩过一次后来在stop脚本里加了一段快速kill和清理逻辑再配合Restarton-failure整体就很稳定了。5. 常见问题、排查思路与实战避坑5.1 高频问题的排查方法速查我把实际运维中遇到最多的问题整理成了一张表每个问题都附上排查方向遇到类似情况可以直接对着查问题现象可能原因排查与解决方向canal server连接MySQL失败账号权限不足、网络不通、server-id冲突检查binlog账号权限、telnet测试3306端口、修改canal.instance.server-id为不冲突值adapter启动后立即退出application.yml格式错误、连接canal超时看adapter日志首行堆栈多半是yaml缩进或连接串错误有增有改但ES里数据不更新同步映射sql的where条件有问题、字段类型不匹配对比adapter日志里的sql执行结果和ES文档内容确认写入是否报错同步延迟越来越大ES写入速度跟不上、bulk批次过大、ES集群负载高调小commitBatch、增加ES节点数、检查是否存在慢查询时间字段差了8小时cannal时区解析和ES不一致检查instance.properties里的时区配置统一到同一时区存储删除事件不同步映射文件缺少delete配置、MySQL没有物理删、binlog格式非row确认binlog-formatROW逻辑删除建议只同步status字段数据重复_id设置不合理导致upsert失效检查映射文件中的_id字段是否对应唯一键避免使用mysql行号当id这个表覆盖了我在几个项目里碰到过的百分之七八十的同步问题来源。每次排查时我第一步永远是看canal server日志和adapter日志的时间戳定位问题发生在“MySQL到canal”还是“canal到ES”缩小范围后对表查找效率高很多。盲目先去改ES映射往往事倍功半。5.2 关于同步延迟和性能优化的经验谈很多人问canal同步到ES到底能做到多快实际受限于两个瓶颈canal从MySQL拉binlog的速度以及ES批量写入的速度。canal本身解析性能非常高瓶颈通常在下游ES。adapter写入ES时有一个commitBatch参数我建议先从3000开始调。commitBatch太小比如500ES请求频繁元数据开销大太大比如1万ES单次bulk压力大容易超时。另外还有个参数叫retry控制写入失败重试次数不要把重试次数设得过高否则ES集群抖动时adapter内部大量堆积重试请求反而拖垮进程。ES端的优化我建议开启refresh间隔调整。默认1秒刷新如果同步量特别大可以在同步时间段把refresh_interval调大到30秒完事再调回来。这个操作不改变最终一致性只是让文档可见的时间延后一些但能显著降低ES段的写入压力。5.3 我最想强调的三个避坑细节第一不要在ES映射里乱用text类型。不少新手把MySQL的varchar字段映射成ES text导致查询时只能用全文检索精确等值匹配全部失效。如果你要的是过滤、聚合、排序使用keyword类型更合适。实在要支持模糊搜索再考虑用text配合keyword子字段。第二canal adapter的sql查询结果一定要保证字段名和ES映射一致。adapter执行完sql后是把返回结果作为整个doc写入ES的。如果字段名大小写不一致或者sql别名和ES映射字段对不上写入后会出现字段缺失或者类型推断错误。第三升级canal版本前一定看一下官方发布的change log。canal的配置文件在版本间有过多次不兼容调整比如从1.1.4到1.1.5adapter的ES连接方式就变过。我早期升级时没注意直接拿旧配置套新版本结果adapter启动后一直连不上ES排查了很久才发现是配置结构的锅。现在我的习惯是每个版本先在一个测试实例上跑通再推生产。写在最后的实战体会跑了这么长时间canal同步MySQL增量数据到ES我的体会是这个方案真正的价值不在于“自动同步”四个字而在于它把数据变更从业务代码里彻底抽离出来了。业务团队不用再关心ES的存在DBA也不用每天晚上跑批数据链路清晰出问题该查哪一层也很明确。最后再分享一个小技巧如果你经常在不同环境部署这套链路建议把canal server、adapter的配置文件和ES索引映射都放进Git仓库用版本号管理。我后来在测试、预发、生产三套环境之间切换靠的就是这套配置即代码的方式环境不一致导致的老大难问题基本消失了。同步链路本身不复杂复杂的是确保它在生产环境里一直稳定运行希望这篇内容能帮你少走几次弯路。