资讯详情

资讯详情

Spring Boot整合Elasticsearch 8.3与RabbitMQ实现MySQL数据同步

简介面向Spring Boot开发者的Elasticsearch 8.3与RabbitMQ集成实战Demo核心解决MySQL数据库数据变化后实时同步至Elasticsearch搜索引擎的典型需求。压缩包共78个文件约90KB包含21个Java源码、17个XML配置、2个YML配置以及Maven包装脚本等Java文件覆盖实体类、Repository接口、消息监听器、服务层与测试逻辑XML/YML文件用于配置依赖、ES连接参数、RabbitMQ队列绑定及项目构建整体目录结构清晰便于导入IDE后按模块阅读和二次改造从依赖配置到索引操作示例一应俱全。项目完整演示了MySQL→RabbitMQ→Spring Boot→Elasticsearch的数据链路并给出批量合并写入、失败重试、日志监控、权限控制等生产环境落地思路同时涵盖ES 8.3与旧版本差异、异步解耦设计等扩展要点。已有209人学习参考适合正在搭建搜索服务或需要异步数据同步方案的中高级开发者直接对照demo即可快速跑通全流程。1. springboot 整合 elasticsearch 8.3 并通过 rabbitmq 同步 mysql 数据库这个 demo 到底要解决什么痛点当一张 MySQL 表到了千万级业务还在用LIKE %关键词%做搜索时接口延时和数据库压力都会肉眼可见地翻车。常见做法是把 MySQL 继续当作唯一数据源把需要检索的字段异步同步一份到 Elasticsearch让 ES 承担搜索MySQL 只负责事务和落地。springboot 整合 elasticsearch 8.3 并通过 rabbitmq 同步 mysql 数据库的 demo就是这条链路的最小可运行版本一个 Spring Boot 服务写业务数据到 MySQL发一条消息到 RabbitMQ再由消费端写入 ES 8.3。这篇文章会从环境选型、客户端配置、生产者消费者实现写到踩坑排查适合准备上 ES、又想先跑通最小闭环的团队。2. 环境与版本选型ES 8.3 从安装到可连接的正确姿势2.1 Windows 启动 Elasticsearch 8.3先把安全认证关掉ES 8.3 和 7.x 最大的差别不只是版本号而是默认开启了安全认证和 TLS。很多人在 Windows 上装完 8.3高高兴兴访问http://localhost:9200结果发现连不上或者拿到一堆证书错误。做 demo 阶段我不建议一上来就跟证书较劲先把elasticsearch.yml里的这两项改掉xpack.security.enabled: false xpack.security.enrollment.enabled: false然后到解压目录执行bin\elasticsearch.bat。启动后看到[o.e.n.Node] started的日志再访问http://localhost:9200能返回带cluster_name的 JSON 就算成功。注意 8.3 对 JDK 版本要求比较高本机没有 JDK 17 的话直接用 zip 包里自带的 JDK 最省事另外 ES 安装路径里不要带空格和中文D:\es\elasticsearch-8.3.3这种路径最稳放在C:\Program Files下会遇到莫名其妙的权限问题。这里还有个容易被忽略的点如果你已经用默认配置启动过一次ES 会在控制台打印elastic用户的初始密码也会生成证书相关文件。想从默认安全模式切到关闭安全改完elasticsearch.yml后一定要重启并且最好把data目录清空否则某些版本仍然会坚持使用已有的安全配置行为上像个黑匣子。我一般会先初始化环境再写业务代码避免把配置问题误判成代码问题。2.2 RabbitMQ 与 MySQL先把待同步的表建出来RabbitMQ 在 Windows 上需要先装 Erlang 再装服务端版本对不上很容易启动失败。我更推荐直接用 Docker 起一个带管理界面的实例一条命令就能把环境准备好docker run -d --name rabbitmq -p 5672:5672 -p 15672:15672 rabbitmq:3.12-management5672是 AMQP 端口给 Spring Boot 连15672是管理台浏览器访问http://localhost:15672默认账号密码都是guest。如果你的服务器上不能跑 Docker那就用本机安装包但记得把端口占用问题优先排查掉。RabbitMQ 启动失败最常见的现象就是端口被占netstat -ano | findstr 5672看一下结果一般能定位到冲突进程。MySQL 这边我用的是 8.0建一个独立库和一张测试表避免污染现有业务数据create database es_demo default character set utf8mb4; use es_demo; create table user_info ( id bigint auto_increment primary key, name varchar(64) not null, age int not null, create_time datetime default current_timestamp ) engineInnoDB;这里有个实战里经常踩的坑Spring Boot 连 MySQL 时连接串一定要加useSSLfalse和serverTimezoneAsia/Shanghai否则后面JdbcTemplate会报 ssl 连接错误时间字段也会差 8 个小时。我在本地跑过的项目里至少见过三次同一个报错每次都是连接串缺参数导致加上就再没出现。2.3 版本矩阵Spring Boot 2.7 elasticsearch-java 8.3.3Spring Boot 3.x 本身当然没问题但它强制 JDK 17很多还在 2.7 上的老项目不会为了一个 demo 就整体升级。为了保证这个 demo 能直接复现我建议用 Spring Boot 2.7.18 elasticsearch-java 8.3.3这也是 ES 8.3 官方客户端里比较稳的一组组合。不要再用RestHighLevelClient。这个客户端在 ES 7.15 就标记废弃到了 8.x 里已经移除继续用只会得到NoClassDefFoundError。正确做法是引入官方 Java API ClientPOM 关键部分如下parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.18/version relativePath/ /parent properties java.version11/java.version elasticsearch.version8.3.3/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-jdbc/artifactId /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId scoperuntime/scope /dependency dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version${elasticsearch.version}/version /dependency dependency groupIdorg.glassfish/groupId artifactIdjakarta.json/artifactId version2.0.1/version /dependency /dependenciesPOM 里最容易漏的是最后那个jakarta.json。Elasticsearch Java API Client 底层 JSON 解析依赖 JSR-374不显式加上会在启动时抛ClassNotFoundException新手经常卡在这。mysql-connector-java的版本交给 Spring Boot 管理即可不用自己写版本号。3. 在 Spring Boot 里接上 ES 8.3客户端配置与索引初始化3.1 application.yml同时接 MySQL、RabbitMQ、Elasticsearch我把三个中间件的配置放在同一个配置文件里方便一眼看出版本和连接方式server: port: 8080 spring: datasource: driver-class-name: com.mysql.cj.jdbc.Driver url: jdbc:mysql://localhost:3306/es_demo?useUnicodetruecharacterEncodingutf8useSSLfalseserverTimezoneAsia/Shanghai username: root password: root rabbitmq: host: localhost port: 5672 username: guest password: guest listener: simple: acknowledge-mode: manual prefetch: 10 elasticsearch: uris: http://localhost:9200 username: elastic password: your_passwordelasticsearch是自定义前缀Spring Boot 不会自动识别需要在自己写的配置类里读取。acknowledge-mode: manual对同步场景很关键后面消费端代码会配合做手动确认prefetch: 10控制消费者预取数量数据量不大时保持默认即可。如果 ES 端已经关闭了安全认证username和password可以随便填因为连接时带的 Basic 头会被忽略如果开启了安全就填实际账号密码。3.2 用官方 Client 创建连接 Bean替代 High Level ClientES 8.3 的常用连接方式有三种Spring Data Elasticsearch、官方 Java API Client、封装后的 RestHighLevelClient但最后一种在 8.x 已经不可用了。我不在这里用 Spring Data因为它对 ES 8.3 的版本绑定比较强升级时容易把 mapper 的坑带出来。官方 Client 更直接配置也透明。Configuration public class EsClientConfig { Value(${elasticsearch.uris}) private String uris; Value(${elasticsearch.username}) private String username; Value(${elasticsearch.password}) private String password; Bean(destroyMethod close) public ElasticsearchClient elasticsearchClient() { RestClient restClient RestClient.builder(HttpHost.create(uris)) .setDefaultHeaders(new Header[]{ new BasicHeader(Authorization, Basic Base64.getEncoder().encodeToString( (username : password).getBytes(StandardCharsets.UTF_8))) }) .build(); ElasticsearchTransport transport new RestClientTransport(restClient, new JacksonJsonpMapper()); return new ElasticsearchClient(transport); } }拆开看这段逻辑RestClient负责 HTTP 层BasicHeader塞进认证头JacksonJsonpMapper负责把 Java 对象和 ES JSON 文档互转。如果后面要调连接超时、socket 超时甚至加自定义拦截器都在RestClient.builder(...)上做二次配置。这个 Bean 的destroyMethodclose很重要不然应用关闭时连接池不会主动释放Spring Boot 优雅停机时可能一直报资源未关闭的警告。3.3 初始化索引mapping 里到底要不要 keywordES 不像 MySQL 有明确的表结构mapping 定义不好检索需求一变就要重建索引。我的习惯是在应用启动时做一个简单的索引初始化没有就创建有就跳过。用CommandLineRunner是最顺手的Component public class EsIndexInitializer implements CommandLineRunner { Autowired private ElasticsearchClient esClient; Override public void run(String... args) throws Exception { String index user_info; boolean exists esClient.indices().exists(e - e.index(index)).value(); if (exists) { return; } esClient.indices().create(c - c .index(index) .mappings(m - m .properties(id, p - p.type(ElasticsearchDataType.Long)) .properties(name, p - p.type(ElasticsearchDataType.Text) .fields(keyword, f - f.type(ElasticsearchDataType.Keyword))) .properties(age, p - p.type(ElasticsearchDataType.Integer)) ) ); } }这里name字段用Text keyword子字段是常见做法全文检索走 text精确匹配和排序走 keyword。ES 8.3 的官方 Client 用 Lambda 风格构建请求看起来比 7.x 的RestHighLevelClient长但类型安全。如果一开始只建了Text后面才发现需要聚合和排序就必须重建索引这是 ES 里典型的没有后悔药的场景所以映射设计宁可一开始多给一个 keyword 子字段。4. 通过 RabbitMQ 把 MySQL 的增删改同步到 ES生产者消费者实现4.1 先选同步链路业务事件消息还是 binlog把 MySQL 数据同步到 ES业界主要两条路一条是业务代码在写库后主动发消息给 MQ另一条是用 Canal、Debezium 这类工具监听 MySQL binlog再转投到 MQ。前者侵入业务但实现简单适合单体应用后者对业务无感但要多维护一套中间件排查链路也更复杂。这个 demo 选的是业务事件消息。原因很简单我们能靠一段 Spring Boot 代码把整条链路跑通不需要额外部署 Canal同时我们也会在事务提交后再发消息尽量保证数据不丢。生产环境如果要求数据库主从复制、批量导入、删除都能无死角同步再考虑引入 binlog 方案也不迟。RabbitMQ 这边的交换机、队列、路由键我先定义成常量后面生产者和消费者共用Configuration public class RabbitConfig { public static final String EXCHANGE es.sync.exchange; public static final String QUEUE es.sync.queue; public static final String ROUTING_KEY user.change; Bean public Queue syncQueue() { return new Queue(QUEUE, true, false, false); } Bean public TopicExchange syncExchange() { return new TopicExchange(EXCHANGE, true, false); } Bean public Binding syncBinding() { return BindingBuilder.bind(syncQueue()).to(syncExchange()).with(ROUTING_KEY); } }队列和交换机都设置成持久化构造函数里的true这样 RabbitMQ 重启后声明不会丢。TopicExchange 对演示来说够用如果以后要按用户、按表拆不同同步任务直接扩展新的 routing key 即可。4.2 生产者MySQL 事务提交后再发 MQ最开始的版本我在事务里直接调用rabbitTemplate.convertAndSend结果发现一旦后面 SQL 回滚消息已经发出去了ES 里出现一条数据库中不存在的脏数据。所以正确顺序是先写 MySQL拿到自增主键等事务afterCommit后再发 MQ 消息。Service public class UserService { private static final ObjectMapper MAPPER new ObjectMapper(); Autowired private JdbcTemplate jdbcTemplate; Autowired private RabbitTemplate rabbitTemplate; Transactional public Long addUser(User user) { KeyHolder keyHolder new GeneratedKeyHolder(); jdbcTemplate.update(con - { PreparedStatement ps con.prepareStatement( insert into user_info(name, age) values(?, ?), PreparedStatement.RETURN_GENERATED_KEYS); ps.setString(1, user.getName()); ps.setInt(2, user.getAge()); return ps; }, keyHolder); Long id keyHolder.getKey().longValue(); TransactionSynchronizationManager.registerSynchronization( new TransactionSynchronization() { Override public void afterCommit() { sendChangeMessage(id, CREATE); } }); return id; } private void sendChangeMessage(Long id, String op) { MapString, Object body new HashMap(); body.put(id, id); body.put(op, op); try { rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE, RabbitConfig.ROUTING_KEY, MAPPER.writeValueAsString(body)); } catch (JsonProcessingException e) { throw new RuntimeException(e); } } }消息体只发主键和操作类型不把整行数据塞进去。这样消费端总是从 MySQL 回查最新数据避免消息中携带的是更新前的旧值。GeneratedKeyHolder的用法是拿自增主键的标准写法MyBatis 项目里改成useGeneratedKeystrue也是一样的逻辑。4.3 消费者从 MySQL 回查数据写入 ES 8.3消费端要做三件事解析消息、按主键回查 MySQL、写或删 ES 文档。这版直接同步完成数据量不大时没问题如果消息量大可以再改成批量聚合后一次写入。Component public class UserSyncConsumer { private static final Logger log LoggerFactory.getLogger(UserSyncConsumer.class); private static final ObjectMapper MAPPER new ObjectMapper(); Autowired private JdbcTemplate jdbcTemplate; Autowired private ElasticsearchClient esClient; RabbitListener(queues RabbitConfig.QUEUE, ackMode MANUAL) public void sync(Message message, Channel channel) throws Exception { long deliveryTag message.getMessageProperties().getDeliveryTag(); try { MapString, Object body MAPPER.readValue( new String(message.getBody(), StandardCharsets.UTF_8), new TypeReferenceMapString, Object() {}); Long id ((Number) body.get(id)).longValue(); String op (String) body.get(op); if (DELETE.equals(op)) { // 删除按 MySQL 主键同步删除 ES 文档 esClient.delete(d - d.index(user_info).id(String.valueOf(id))); } else { // 更新/新增回查最新数据再写 ES User user jdbcTemplate.queryForObject( select id, name, age from user_info where id ?, (rs, rowNum) - new User(rs.getLong(id), rs.getString(name), rs.getInt(age)), id); if (user ! null) { esClient.index(i - i.index(user_info) .id(String.valueOf(user.getId())) .document(user)); } } // 处理成功才确认 channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(sync to es failed, message tag {}, deliveryTag, e); channel.basicNack(deliveryTag, false, true); } } }这里最关键的是id(String.valueOf(user.getId()))把 MySQL 主键直接对应到 ES 文档的_id。如果不指定ES 会自动生成随机 id同一行数据多次更新就会产生多条重复文档这是同步链路里最隐蔽的数据质量事故。DELETE 分支即使 ES 里没有对应文档删除也会返回not_found不会抛异常所以不用担心重复消费。4.4 手动确认、重试与幂等同步链路的最后一道锁消费者里用了ackMode MANUAL对应配置文件的acknowledge-mode: manual。手动 ack 的含义很清楚处理成功才basicAck失败就basicNack。但basicNack的第三个参数requeue要特别小心我见过一条无法解析的坏消息被无限重新投递把队列后面的正常消息全都堵死。更稳妥的退避策略是把requeue设为false配合死信队列保存失败消息再通过定时任务人工补偿。这个在 demo 里不展开但要知道方向。重试的幂等性这里天然满足因为消费者每次都按 MySQL 主键回查最新数据同一个 CREATE 消息消费两次ES 文档会覆盖成同一份不会产生双份数据。5. 避坑记录ES 8.3 和 RabbitMQ 同步链路的 5 个血泪教训下面五条是我跑这个 demo 时真正遇到过的坑每一条都按现象、原因、解决整理照着排查能少走弯路。5.1 现象启动报NoClassDefFoundError: RestHighLevelClient原因项目里还在用 7.x 时代的RestHighLevelClientES 8.3 已经把相关类移除。Spring Boot 自带的spring-boot-starter-data-elasticsearch在某些版本组合下也会把旧客户端带上导致编译通过、运行时报类找不到。解决放弃 high-level client统一用co.elastic.clients:elasticsearch-java:8.3.3并保证客户端版本和 ES 服务器大版本一致。8.3.3 的客户端不去连 8.0 或 8.2也不要指望 7.x 客户端兼容 8.3ES 官方在 8.x 上对版本匹配要求很严格。5.2 现象RabbitMQ 第一次发消息就报reply-code404, reply-textNOT_FOUND原因生产者把消息发到交换机时交换机或队列还没有在当前连接中存在。如果只声明队列但没声明交换机或者路由键对不上RabbitMQ 会直接关闭信道日志里出现clean channel shutdown; protocol method: #methodchannel.close(reply-code404, reply-textNOT_FOUND...)这一串看起来像乱码的东西。解决在 Spring Boot 里把Queue、TopicExchange、Binding三个 Bean 都定义好应用启动时自动声明。然后在管理台http://localhost:15672里看 Exchange 和 Queues 是否存在再用 Publish Message 手工发一条测试消息确认路由能到队列。这个消息排查顺序很有效能快速区分是声明问题还是消费端问题。5.3 现象ES 8.3 已经启动但客户端用 HTTP 一直连接失败原因ES 8.3 默认开启安全TLS 也是默认启用的访问地址是https://localhost:9200而不是http://localhost:9200。用 HTTP 连过去要么收到证书错误要么连接被重置。如果你在启动日志里看到elastic用户的初始密码说明安全功能已经生效证书的坑也一起埋下了。解决本地 demo 阶段在elasticsearch.yml里显式关闭xpack.security.enabled和xpack.security.enrollment.enabled重启后再用 HTTP 连接。要保留安全功能的话就用 HTTPS并把 CA 证书加载到TrustStore再在RestClient上配置SSLContext。我自己的经验是不要一边开着安全一边用 HTTP 调试这个问题最容易让人误判成 Spring Boot 配置写错。5.4 现象消费者里明明写basicAck但消息还是被自动确认了原因注解上ackMode MANUAL没生效或者配置文件里spring.rabbitmq.listener.simple.acknowledge-mode写的是小写manual但容器工厂没有读到。还有一种更隐蔽的情况消费者方法抛异常后没有basicNack消息被容器自动拒绝并重新入队表现成队列里的 Unacked 一直不变你以为 ack 没生效。解决先确认配置路径是listener.simple.acknowledge-mode不是listener.direct。在RabbitListener上显式声明ackMode MANUAL并确保消费端代码里每个分支都有basicAck或basicNack。不要把手动确认和自动确认混用一个方法里全部手动整体才可控。5.5 现象同步完成后去 ES 查询经常查不到刚写入的文档原因ES 写入是近实时的默认 refresh interval 是 1 秒写入后立即查询大概率落在 refresh 之前所以查不到。这个不是同步失败也不是数据丢了而是索引还没刷新。解决调试阶段可以在查询时加?refresh参数或者调用客户端时设置refresh(Refresh.WaitFor)强制等待写入完成再返回。正式同步链路不要这么干会让写入吞吐明显下降。初始全量导入时可以把index.refresh_interval设成-1导入完成再改回1s能省不少时间。这个参数在同步场景里非常有用。6. 从 demo 到能用验证链路和三个值得调的参数链路写完了第一件事不是写更多代码而是验证整条链路。启动 Spring Boot 后先调接口插入一条用户curl -X POST http://localhost:8080/user \ -H Content-Type: application/json \ -d {name:张三,age:25}返回里带上新生成的主键后手动等 1 秒再查 EScurl http://localhost:9200/user_info/_search?qname:张三pretty能查到文档说明 MySQL 写入、RabbitMQ 消息、消费端回查、ES 索引四个环节都正常。接下来把 RabbitMQ 管理台打开看es.sync.queue的Ready和Unacked数量Ready应该很快归零Unacked如果一直上涨优先怀疑消费端处理慢或抛异常。这条验证路径比直接写单元测试更贴近生产因为消息中间件的状态只有真实跑起来才能反映。生产化之前我建议把三个参数单独调一遍。第一个是 RabbitMQ 的prefetch默认 250 对 ES 写入偏激进同步任务里 10 到 30 更稳避免消费太快把 ES bulker 打满第二个是 ES 客户端的批量接口demo 里是一条一条index数据量上来后要改成BulkRequest打包写入一次提交几百条吞吐能差一个量级第三个是索引的refresh_interval测试环境保持默认初始全量同步时调成-1跑完再调回来同步时间能缩短不少。我做这个 demo 时踩得最多的是版本和安全配置最后留下一个习惯每次拿到新的 ES 大版本先看 release notes 里客户端和服务器版本的对应关系再在空分支上跑通最小链路而不是直接往老项目里塞依赖。ES 8.3 的官方客户端和 RabbitMQ 这套组合只要版本对齐、事务提交后再发消息、消费端手动确认并幂等写入数据同步这条链路是能睡个安稳觉的。希望帮到你。本文还有配套的精品资源点击获取
觉得有用,分享给同行:

为您的企业打造数字门面

稳重轻奢商务风格,端正雅致视觉,长效耐看不易过时。

立即咨询 →