拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Spring Boot整合Elasticsearch 8.3:RabbitMQ同步MySQL数据全链路实战

简介面向需要将 MySQL 中的数据实时同步至 Elasticsearch 搜索引擎的 Java 开发者这套基于 Spring Boot 整合 Elasticsearch 8.3 的 Demo借助 RabbitMQ 消息中间件设计了一套完整且可运行的异步同步机制。资源包内合计 78 个文件整体大小约 90KB包括用于业务逻辑的 Java 源码、定义依赖与持久层映射的 XML 配置、描述环境参数的 YAML 文件以及保障通信安全的 SSL 证书工程结构紧凑方便直接导入开发环境学习验证。目前已有 208 人学习下载。这套工程从依赖引入开始逐步演示了 Spring Data Elasticsearch 的 Repository 接口操作、RabbitMQ 消息生产者与监听器的编写、基于 TransactionalEventListener 捕获数据库变更事件再通过消息驱动完成 Elasticsearch 索引更新的完整数据流。此外还给出了批量写入、异常重试、监控安全等优化思路可帮助读者快速搭建一套高可用的 MySQL 向 Elasticsearch 同步的基础工程并迁移至电商搜索、日志分析等实际业务场景。 把“springboot整合elasticsearch8.3并通过rabbitMq同步mysql数据库的demo”这个需求扔进搜索引擎搜出来的教程十有七八还停留在High Level REST Client的写法。ES 8.x把老客户端移除之后大量旧代码直接报废网上资料又是七零八落拼凑出来的照着敲完根本跑不通。这篇文章就围绕一条完整可跑的链路来讲Spring Boot 2.7 Elasticsearch 8.3 RabbitMQ MySQL用MQ做数据同步的桥梁把MySQL里的业务数据增量推到ES索引里。我会把环境版本、客户端初始化、生产者消费者、增量查询、幂等更新和踩坑实录全部串起来适合正在搭第一个ES同步demo、又被各种版本兼容问题折磨的人参考。1. 为什么不是Canal也不是Logstash而是MQ手动同步1.1 小数据量下的控制欲定时增量查询完全够用接触过ES同步的人应该都听过Canal和Logstash但真正动手做demo时这两个方案反而容易把人带偏。Canal需要额外部署一个Java进程去伪装成MySQL的从库通过解析binlog拿到变更记录这要求MySQL开启binlog并设置成ROW模式还要给Canal单独开一个账号授权。Logstash虽然支持JDBC输入插件直接轮询MySQL但它的同步频率、增量字段处理、删除同步都是配置项驱动的一旦业务逻辑复杂配置文件的维护成本比写代码高得多。MQ手动同步的思路更直接业务代码里操作完MySQL把变更数据封装成消息丢到RabbitMQ消费者收到消息后调用ES的Java API Client写入索引。整个过程不需要额外部署中间件不需要理解binlog的解析规则所有逻辑都在项目里可见、可断点、可改。对单体应用和中小体量的数据同步场景这种“裸写”的方式反而比引入一堆重型组件更可控排查问题时也只需要盯住一条消息的流转即可。1.2 三条路线的真实对比不同方案的核心差异在于谁负责感知MySQL的变更以及变更数据以什么形态进入ES。我用一个表格把这三种主流方案的定位说清楚方案变更感知方式额外依赖适用场景主要痛点CanalMySQL binlog监听需部署Canal服务端高并发、强一致、大规模数据迁移运维成本高Canal版本与MySQL 8.0的认证插件兼容性容易出问题LogstashJDBC定时轮询需部署Logstash日志类、离线批量同步实时性受轮询间隔限制复杂的业务映射需要写大量filterRabbitMQ手动同步业务代码显式发送消息仅依赖MQ本来就该有业务系统内部数据同步同步逻辑与业务强相关需要自己处理消息可靠性和幂等看到这里你应该明白了Canal适合“不动业务代码”的旁路同步场景Logstash适合日志和离线管道而MQ手动同步适合业务系统内部的精准控制。三者不是替代关系是共存关系。demo阶段用MQ手动同步能最快把整条链路吃透之后再切换到其他方案也有扎实的基础。2. 环境配置里最容易翻车的三个地方2.1 版本矩阵Spring Boot、ES 8.3、JDK、RabbitMQES 8.x的客户端包做了一次大重构把原本的High Level REST Client整个废弃换成了Elasticsearch Java API Client包名也变成了co.elastic.clients。这个改动直接影响pom依赖的引入方式和代码写法。另一个问题是JDK版本ES 8.3服务端强制要求JDK 17而Spring Boot 2.7系列默认兼容JDK 8和11如果本机只装了JDK 8ES服务端根本起不来。我实测下来这套组合最稳组件推荐版本备注JDK17ES 8.3服务端强依赖Spring Boot 2.7也兼容Spring Boot2.7.x3.x把javax迁移到jakartaES客户端没跟上时容易踩雷Elasticsearch8.3.0不要用8.0/8.1部分API还不够稳定RabbitMQ3.10用Docker跑最省事MySQL8.05.7也能跑通但8.0更贴近生产环境Spring Boot父pom里配置ES版本有个坑ES的Java API Client不会跟着Spring Boot的依赖管理自动对齐版本需要显式声明。下面这段pom.xml是跑通的配置parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.4/version relativePath/ /parent properties elasticsearch.version8.3.0/elasticsearch.version /properties dependencies dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-elasticsearch/artifactId /dependency dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version${elasticsearch.version}/version /dependency dependency groupIdcom.baomidou/groupId artifactIdmybatis-plus-boot-starter/artifactId version3.5.3/version /dependency /dependencies注意spring-boot-starter-data-elasticsearch这个依赖Spring Data ES 4.4版本已经能适配ES 8.3的客户端但它在ElasticsearchRestTemplate里的查询写法仍然偏老所以我在实践里选择直接用原生Java API Client不走Spring Data封装层减少一层“魔法”也更方便排查问题。2.2 ES 8.3客户端初始化的“新写法”ES 8.3初始化客户端和7.x完全不同。7.x时代要用RestHighLevelClient8.3必须在RestClientTransport的基础上构建ElasticsearchClient并且序列化器要指定JSON映射器。我用Jackson这样返回的实体可以直接和数据库实体类互转。Configuration public class ElasticsearchConfig { Bean public ElasticsearchClient elasticsearchClient() { RestClientBuilder builder RestClient.builder( new HttpHost(localhost, 9200, http) ); RestClient restClient builder.build(); RestClientTransport transport new RestClientTransport( restClient, new JacksonJsonpMapper() ); return new ElasticsearchClient(transport); } }一个小细节ES 8.3默认开启了安全认证本地测试如果不想配用户名密码在Docker启动时加上-e xpack.security.enabledfalse可以免认证访问。但一旦关闭9200端口就成了完全开放状态只建议在本机或内网测试环境这么搞千万别用到公网服务器上。2.3 索引mapping与中文分词插件的匹配索引mapping必须在写入数据前定义好否则ES会按默认策略动态映射。默认的standard分词器对中文是按字切分的搜“技术”搜不出“互联网技术”里的“技术”这个词中文检索必须配合IK分词插件。ES 8.3.0对应IK分词器的版本是8.3.0从medcl/elasticsearch-analysis-ik的release页面下载对应的zip包解压到ES安装目录的plugins/analysis-ik文件夹下重启ES即可。创建索引时指定IK分词器PUT /user_index { mappings: { properties: { id: {type: integer}, name: { type: text, analyzer: ik_max_word, search_analyzer: ik_smart }, email: {type: keyword}, createTime: {type: date} } } }这里有个经验谈ik_max_word和ik_smart的分词粒度不同ik_max_word更细适合索引阶段把文本拆分成尽可能多的词ik_smart更粗适合搜索阶段匹配更精准的结果。如果两个阶段用同一个分词器会导致查“中华人民共和国”时把“中华”也匹配出来精确度下降。3. 数据同步链路的代码实现3.1 MySQL侧基于update_time做增量查询同步的第一步是从MySQL查出变更数据。demo阶段最稳的做法是加一个update_time字段每次更新记录都自动刷新这个时间戳然后消费者定时查询“最近一段时间内有修改”的数据。SELECT id, name, email, create_time, update_time FROM user WHERE update_time #{lastSyncTime} ORDER BY update_time ASC这个方案简单、不侵入业务代码唯一的缺点是删除的数据无法被发现。要处理删除就得在业务删除代码里同步发一条删除消息或者在表里加逻辑删除标志位。实际项目中我见过太多只做增量不同步删除的案例导致ES里沉淀一堆脏数据所以从这个demo开始就要把删除链路想清楚。结合MQ的驱动方式我的做法是生产者服务在MySQL的事务提交之后把变更的实体对象序列化为JSON发到RabbitMQ。这里有个细节——必须在事务提交后再发消息否则事务回滚了消息却发出去了消费端拿到一条根本不存在的数据。3.2 RabbitMQ侧队列、交换机、路由键设计RabbitMQ的队列设计直接影响消息分发的灵活性。我可以为不同业务表建不同的队列也可以把所有变更消息都丢进一个队列消费端根据消息里的type字段做路由。demo阶段我倾向于后者因为代码结构更简单模块边界清晰还为以后增加新的同步表留了扩展位。RabbitMQ的配置类Configuration public class RabbitConfig { public static final String EXCHANGE_NAME es.sync.exchange; public static final String QUEUE_NAME es.sync.queue; public static final String ROUTING_KEY es.sync.user; Bean public DirectExchange esSyncExchange() { return new DirectExchange(EXCHANGE_NAME, true, false); } Bean public Queue esSyncQueue() { return QueueBuilder.durable(QUEUE_NAME).build(); } Bean public Binding esSyncBinding() { return BindingBuilder.bind(esSyncQueue()) .to(esSyncExchange()) .with(ROUTING_KEY); } }消息体我建议用JSON而不是Java对象序列化。原因很实际Java原生序列化把类结构耦合在消息里一旦消费者升级改了实体字段老消息反序列化直接报错JSON是纯文本消费者可以灵活容错还能在RabbitMQ管理后台直接看到消息内容排查问题极其方便。3.3 ES侧Java API Client写入索引消费者拿到消息后把JSON解析成实体对象调用ES客户端写入。这里核心是index方法的Lambda表达式写法和7.x的IndexRequest完全是两回事Component public class UserSyncConsumer { Autowired private ElasticsearchClient esClient; RabbitListener(queues RabbitConfig.QUEUE_NAME) RabbitHandler public void handleMessage(String message) throws Exception { UserEntity user JSON.parseObject(message, UserEntity.class); IndexResponse response esClient.index(i - i .index(user_index) .id(String.valueOf(user.getId())) .document(user) ); if (response.result() Result.Created || response.result() Result.Updated) { log.info(同步成功id: {}, user.getId()); } } }.id()必须显式指定为MySQL里的主键这样ES的主键就和MySQL一致重复投递消息时用同一个id覆盖写入天然实现了幂等更新。如果漏了.id()ES会对每次写入生成随机id重复消费一条消息就会产生多条重复文档这个坑非常隐蔽。搜索时用match查询验证数据是否同步成功SearchResponseUserEntity search esClient.search(s - s .index(user_index) .query(q - q.match(m - m .field(name) .query(张三) )), UserEntity.class );4. 几个必须想清楚的工程问题4.1 消息丢失与重复消费让更新操作天然幂等用MQ做同步最怕的就是消息丢了或者消费者处理到一半宕机。RabbitMQ的消息可靠性有三个层级生产者确认机制、队列持久化、消费者手动确认。三者缺一不可。生产者开启publisher-confirm-type: correlated发送后异步回调确认消息是否到达交换机队列声明时设置durable(true)消费者用manual模式的ack处理只有ES写入成功后才会返回basicAck。如果ES写入失败消息重新入队重试超过次数后进入死信队列。”手动确认“和“幂等更新”是一对好搭档就算消费者在处理完ES写入后、还没来得及发送ack时宕机消息会重新投递。但因为ES写入的id就是MySQL的主键重复执行覆盖更新并不会产生副作用这个组合才能成立。4.2 MySQL库表结构变更后ES映射怎么办这是生产环境一定会遇到的问题MySQL表加了字段ES的索引映射没跟上消费者反序列化后写入时新的字段没有对应的ES字段映射要么写不进去要么动态映射出了意外的类型。要根治这个问题需要一套映射管理机制MySQL表结构每次变更都要同步更新ES索引的mapping。我的做法是在项目里维护一个索引初始化脚本配合ES的indices.create幂等性在启动时检查并自动创建缺失的索引。这个脚本必须是可重复执行的因为ES不允许修改已有字段的类型只能新增字段。如果业务实在需要改类型只能重建索引再reindex这是ES的硬限制早点知道早点设计好。4.3 从demo到生产的差距在哪这个demo跑通之后距离生产还有很长的路要走。生产者发送消息前要做消息体校验避免把脏数据同步过去消费端要增加重试退避机制防止ES短时不可用时消息风暴式重投还要有监控大盘把消息积压数、同步失败数、ES写入耗时全部可视化。更深一层的问题是“双写一致性”。业务在MySQL上做了修改消息发出去了但是ES写入前被其他操作打断这个短暂的时间窗口里MySQL和ES的数据是不一致的。对于搜索场景这个窗口通常可以接受但如果业务要求强一致就需要引入本地消息表把消息和业务数据放在同一个数据库事务里通过定时任务兜底补偿。这套机制是很多秒杀、支付系统的标准做法放在ES同步链路里同样适用。5. 踩坑实录与调试经验5.1 连不上的排查链路ES服务端起不来或者客户端连不上是所有人都会遇到的第一道坎。我遇到过最常见的几个报错报错信息原因解决Connection refusedES服务没起或者端口不对确认9200为HTTP端口9300为Transport端口已经废弃了missing authentication credentialsES 8默认开启了安全认证测试环境启动时加xpack.security.enabledfalsejavax.net.ssl.SSLHandshakeException客户端走http但是ES期望https确认RestClient.builder()里的协议是http还是httpsUnsupported class file major version 61.0JDK版本低于17把项目JDK切换到17排查顺序我有一个标准流程先curl localhost:9200看是否返回JSON再用客户端的info()方法验证Java代码能连通最后才去查索引和数据。这样逐层排除比在代码里到处打日志高效得多。5.2 ik分词不生效的根因很多人按网上的教程装好了IK分词器却发现搜索“华为”还是匹配不到“华为技术有限公司”。这里有个很容易忽略的点分词器是在索引创建时绑定到字段上的修改IK插件或者更新mapping配置对已存在的索引不生效。具体来说如果你先创建了索引用的是默认standard分词器再安装IK插件那么旧索引里的字段仍然是standard分词器必须删除索引重新创建或者用reindex API把数据迁移到新索引。我在本地测试时习惯把索引生命周期脚本和IK插件安装放在同一批初始化动作里避免顺序错乱。验证分词效果有专门的接口不用写代码POST /user_index/_analyze { field: name, text: 华为技术有限公司 }返回的tokens数组里如果出现了“华为”“技术”“有限公司”这些词说明IK正常工作。这个方法在调试搜索相关问题时非常快强烈建议掌握。5.3 写进去了搜不到ES写入成功后立刻搜索经常返回空结果。这不是代码错了是ES的refresh机制导致的。ES写入的数据先进内存缓冲默认每秒做一次refresh把缓冲区的数据写入新的segment并打开它之后搜索才能命中。所以刚写入的文档在极端情况下会有接近1秒的搜索延迟。demo阶段为了立竿见影可以在写入后手动刷新索引esClient.indices().refresh(r - r.index(user_index));生产环境不建议这么做高频刷新会严重拉低写入性能。更好的做法是确认搜索场景能接受1秒内的延迟如果确实需要近实时把refresh_interval从默认的1s调低到500ms但提升了实时性的同时也要承担更多的segment合并开销。关于调试还有一个经验RabbitMQ管理后台的队列页面可以直观看到消息的入队、出队和堆积状态。如果消费者一直不消费先检查队列的consumer count是否为1再看消费者日志里有没有反序列化异常最后用管理后台的“Get message”功能手动拉一条消息看内容格式三步就能定位大多数同步链路问题。最后再分享一个技巧消费者监听队列时把日志级别调到DEBUG跑一次全链路能同时看到RabbitMQ的消息族、ES请求的DSL对理解整个数据流向帮助巨大比单步debug高效得多。我第一次把这套链路完整跑通时最深的体会是ES 8.3的Java API Client虽然改了API风格但核心思路没变无非还是“建立客户端、指定索引、执行写入、发起查询”四件事把MQ当作中间的缓冲管道这条链路并没有想象中那么复杂。本文还有配套的精品资源点击获取
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门