ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SpringBoot整合ES8.3并用RabbitMQ同步MySQL数据

SpringBoot整合ES8.3并用RabbitMQ同步MySQL数据 简介这是面向Spring Boot开发者的Elasticsearch 8.3集成与数据同步实战Demo主要解决关系型数据库到搜索引擎的实时同步问题。资源通过RabbitMQ消息队列打通MySQL与Elasticsearch涵盖Spring Data Elasticsearch配置、Repository接口定义、消息监听与发送、数据库变更事件处理等完整链路适合有一定Spring Boot基础、需要构建搜索服务的后端工程师。压缩包共78个文件包含21个Java源码、21个编译后class、17个XML配置、2个YML环境配置及Maven脚本等目录中源码与构建产物分离清晰便于直接对照学习和调试整体仅90KB。目前已有209人学习轻量易用。从中可掌握微服务架构下MySQL到Elasticsearch的实时同步思路、批量写入与错误处理技巧以及项目配置和测试监控的落地写法可直接迁移到电商搜索、日志分析等实际场景中。1. 这条 Demo 在解决什么问题MySQL 负责存Elasticsearch 负责搜「springboot整合elasticsearch8.3并通过rabbitMq同步mysql数据库的demo」这个标题核心就一句话业务数据在 MySQL搜索需求在 Elasticsearch 8.3中间的管道用 RabbitMQ 异步打通。最常见的切入场景是商品表或者订单表到了几十万行MySQL 的LIKE %keyword%查询把接口拖到几秒甚至超时于是把 ES 拉进来扛全文检索。但 ES 不能当数据主库它只是 MySQL 的一份索引副本主数据一旦变了副本必须跟着变——这就是整条同步链存在的意义。这篇文章面向已经能写 Spring Boot 接口、但没深入碰过 Elasticsearch 和 RabbitMQ 的开发者按能跑通的顺序讲清楚ES 8.3 客户端怎么接、RabbitMQ 消息怎么发怎么收、消费端怎么把 MySQL 数据写进 ES、以及最容易翻车的几个地方。2. Elasticsearch 8.3 与 Spring Boot 整合为什么不用 Spring Data用官方 Java API Client2.1 版本适配是最容易掉进去的坑RestHighLevelClient 在 8.x 里已经废弃用 Spring Boot 整合 Elasticsearch第一反应通常是引入spring-boot-starter-data-elasticsearch但这里有个历史包袱Spring Data Elasticsearch 在 Spring Boot 2.7.x 里默认对接的是 ES 7.x底层依赖RestHighLevelClient。到了 ES 8.x官方把RestHighLevelClient标记为废弃主推新的Elasticsearch Java API Client也就是elasticsearch-java这个库。强行用 Spring Data 去连 ES 8.3最常见的报错是版本冲突或者序列化方式不兼容最后 debug 半天还是得换掉。所以这条链路的常见做法是不走 Spring Data直接用官方客户端ElasticsearchClient它底层走 HTTP和 Spring Boot 没有版本强绑定相当于在工程里引入一个普通的 HTTP 客户端 Bean。2.2 最小依赖与连接配置pom.xml、application.yml 和配置类一次到位先建一个空白的 Spring Boot 工程2.7.x 即可引入两个依赖版本必须对齐到 8.3.0dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version8.3.0/version /dependency dependency groupIdorg.elasticsearch.client/groupId artifactIdelasticsearch-rest-client/artifactId version8.3.0/version /dependency第一行是 ES 8.3 的官方 Java API Client负责把 Java 对象映射成 JSON 请求第二行是底层 HTTP 连接器真正和 ES 节点通信。注意这两个版本号必须一致如果你把elasticsearch-java升到 8.10、elasticsearch-rest-client还停在 8.3运行时会直接甩NoSuchMethodError属于最典型的版本对齐问题。ES 8.3 要求 JDK 17 起步Spring Boot 2.7 跑在 JDK 17 上没有问题但如果你本地还是 JDK 8这一步就得先升级。接着在application.yml里放连接信息。这里有个细节spring.elasticsearch.uris这套配置是 Spring Data Elasticsearch 的官方客户端不认所以更稳的做法是自己写一个配置类通过Value读配置es: uris: http://localhost:9200 username: elastic password: 123456ES 8.3 默认开启了 xpack 安全认证直接用RestClient.builder(new HttpHost(localhost, 9200))去连会收到401 Unauthorized。有两种解法要么本地测试时在 docker 启动命令里关掉安全认证要么在客户端里带上 Basic Auth。下面这个配置类采用带认证的方式Configuration public class EsClientConfig { Value(${es.uris}) private String uris; Value(${es.username}) private String username; Value(${es.password}) private String password; Bean public ElasticsearchClient elasticsearchClient() { HttpHost host HttpHost.create(uris); BasicCredentialsProvider credentialsProvider new BasicCredentialsProvider(); credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password)); RestClient restClient RestClient.builder(host) .setHttpClientConfigCallback(hcb - hcb.setDefaultCredentialsProvider(credentialsProvider)) .build(); ElasticsearchTransport transport new RestClientTransport(restClient, new JacksonJsonpMapper()); return new ElasticsearchClient(transport); } }逻辑说明HttpHost.create(uris)会把http://localhost:9200解析成 host 和端口BasicCredentialsProvider是 Apache HttpClient 的认证机制ES 8.3 开启安全认证后所有请求头里必须带Authorization: Basic ...这个回调会自动附加JacksonJsonpMapper负责把 Java 对象序列化成 ES 能读懂的 JSON同时把 ES 返回的 JSON 反序列化成 Java 对象。如果你在本地用 docker 关掉了安全认证把UsernamePasswordCredentials这段删掉即可。2.3 建索引、写入、查询跑通 ES 8.3 的最小可用代码客户端 Bean 配好之后先写一个服务类验证连通性顺便把索引映射建好。这一步非常重要ES 的字段类型映射mapping一旦索引创建后就不好改了所以建索引时就要把 MySQL 字段和 ES 字段的类型对应关系定清楚。Service public class ProductIndexService { private final ElasticsearchClient client; public ProductIndexService(ElasticsearchClient client) { this.client client; } public void createIndexIfNotExists(String indexName) throws IOException { boolean exists client.indices().exists(e - e.index(indexName)).value(); if (!exists) { client.indices().create(c - c.index(indexName) .mappings(m - m .properties(id, p - p.keyword()) .properties(title, p - p.text()) .properties(price, p - p.double_()) .properties(createTime, p - p.date() .format(yyyy-MM-dd HH:mm:ss)) ) ); } } public void saveDoc(String indexName, String id, Object doc) throws IOException { client.index(i - i.index(indexName) .id(id) .document(doc)); } public void deleteDoc(String indexName, String id) throws IOException { client.delete(d - d.index(indexName).id(id)); } }逻辑说明client.indices().exists返回一个布尔值避免重复创建索引报错properties里面逐个声明字段类型id用 keyword因为主键只做精确匹配不需要分词title用 text因为商品名称需要分词检索price用 double对应 MySQL 的decimal转过来createTime用 date 并指定格式这样 ES 才能识别2024-01-01 12:00:00这种字符串并进行范围查询。saveDoc里显式传了id这是为后续 RabbitMQ 消息重复投递做的准备同一个id重复写入 ES 会直接覆盖天然幂等。写一个测试接口快速验证RestController public class IndexTestController { private final ProductIndexService indexService; public IndexTestController(ProductIndexService indexService) { this.indexService indexService; } PostMapping(/index/test) public String test() throws IOException { indexService.createIndexIfNotExists(product); MapString, Object doc new HashMap(); doc.put(id, 1); doc.put(title, 机械键盘 红轴); doc.put(price, 399.00); doc.put(createTime, 2024-01-01 12:00:00); indexService.saveDoc(product, 1, doc); return ok; } }这一步做完curl http://localhost:9200/product/_search能看到刚写入的数据ES 8.3 和 Spring Boot 的整合就算通了。注意启动参数里的内存问题ES 8.3 默认堆内存 1GB机器内存小的话加ES_JAVA_OPTS-Xms512m -Xmx512m再启动不然 ES 进程起不来这是 Windows 上启动 Elasticsearch 最常见的问题。3. RabbitMQ 消息链路从 MySQL 变更到队列消息的完整设计3.1 交换机类型选 TopicRoutingKey 按业务语义命名RabbitMQ 的角色是中间管道MySQL 的数据变更被包装成一条消息扔进队列消费者在另一端接收并回查 MySQL 后写入 ES。交换机类型我一般选 Topic不用 Fanout也不用 Direct。原因很简单Fanout 会把消息广播给所有绑定的队列将来加了别的消费者比如同步到缓存、做数据统计不想让它们都收到全量数据Direct 又太死板只能精确匹配 routing key。Topic 允许mysql.es.sync和mysql.es.retry这类带点号的路由规则同一个交换机后续扩展死信队列、重试队列时不用换拓扑。3.2 声明交换机、队列、绑定关系一份配置类全部搞定Configuration public class RabbitSyncConfig { public static final String SYNC_EXCHANGE exchange.mysql.es; public static final String SYNC_QUEUE queue.mysql.es.sync; public static final String SYNC_ROUTING_KEY mysql.es.sync; Bean public TopicExchange syncExchange() { return new TopicExchange(SYNC_EXCHANGE, true, false); } Bean public Queue syncQueue() { return new Queue(SYNC_QUEUE, true); } Bean public Binding syncBinding() { return BindingBuilder.bind(syncQueue()) .to(syncExchange()) .with(SYNC_ROUTING_KEY); } }逻辑说明new TopicExchange(SYNC_EXCHANGE, true, false)三个参数分别是名字、是否持久化、是否自动删除生产环境持久化必须为 true否则 RabbitMQ 重启后交换机就没了。Queue同样设置持久化防止消息进来之后 broker 重启导致队列数据丢失。Binding把队列和交换机通过 routing key 绑起来消息带着mysql.es.sync这个 key 发进来队列就能收到。这段配置是静态拓扑Spring Boot 启动时自动声明不会重复创建。消息体不建议直接塞整行 MySQL 数据而是只放必要信息配合消费者端回查数据库public class SyncMessage { private String op; // INSERT / UPDATE / DELETE private String tableName; // 表名 private Long businessId; // MySQL 主键 ID private Long timestamp; // 操作时间 // 无参构造、getter/setter 省略 }这样做有三个好处一是消息体积小RabbitMQ 的吞吐压力小二是消费者总能拿到 MySQL 里的最新数据不会因为消息在链路里排队而把旧值写进 ES三是敏感字段不会出现在 MQ 消息里排查问题时看到的是主键而不是整行用户信息。代价是消费者必须查一次 MySQL这对同步场景来说可以接受。3.3 发送端事务提交后再发消息别在事务里直接发发送端最关键的坑是事务边界。如果在Transactional方法里直接调用rabbitTemplate.convertAndSend消息发出去了但 MySQL 事务还没提交消费者如果动作够快回查到的是旧数据等事务回滚时消息就彻底变成一条脏消息。常见的正确做法是注册事务同步器在afterCommit回调里发消息Service public class ProductService { Autowired private ProductMapper productMapper; Autowired private RabbitTemplate rabbitTemplate; Transactional public void save(Product product) { productMapper.insert(product); TransactionSynchronizationManager.registerSynchronization(new TransactionSynchronization() { Override public void afterCommit() { SyncMessage msg new SyncMessage(); msg.setOp(INSERT); msg.setTableName(product); msg.setBusinessId(product.getId()); msg.setTimestamp(System.currentTimeMillis()); rabbitTemplate.convertAndSend( RabbitSyncConfig.SYNC_EXCHANGE, RabbitSyncConfig.SYNC_ROUTING_KEY, msg); } }); } }逻辑说明TransactionSynchronizationManager.registerSynchronization是 Spring 事务抽象的扩展点afterCommit方法只有在事务真正提交成功后才会执行。这样保证了「MySQL 有数据」和「MQ 有消息」两个动作在时间顺序上是安全的。convertAndSend的三个参数分别是交换机名、routing key、消息体。这里还差一步需要把消息序列化方式从 JDK 自带的序列化改成 JSON否则消息在 RabbitMQ 管理后台里是一堆二进制乱码跨语言消费者也读不了Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); }这段配置加在任意Configuration类里即可Spring Boot 会自动把它注入RabbitTemplate。改完之后RabbitMQ 控制台里能看到可读的 JSON 消息排查问题会舒服很多。发送端的异步确认建议开启spring.rabbitmq.publisher-confirm-type: correlated这样convertAndSend之后可以通过CorrelationData回调确认 broker 是否真正收到了消息——但这是优化项先把链路跑通再说。4. 消费者落地回查 MySQL写入 Elasticsearch处理三种操作类型4.1 RabbitListener 消费按消息里的主键回查数据库消费者是这条链路的终点职责很单一收到消息、按主键回查 MySQL、把数据写入 ES。用一个Component加RabbitListener即可Component public class ProductSyncConsumer { Autowired private ProductMapper productMapper; Autowired private ProductIndexService indexService; RabbitListener(queues RabbitSyncConfig.SYNC_QUEUE) public void onMessage(SyncMessage message) { if (DELETE.equals(message.getOp())) { indexService.deleteDoc(product, String.valueOf(message.getBusinessId())); return; } Product product productMapper.selectById(message.getBusinessId()); if (product null) { // 数据不存在说明可能已被删除或事务未提交直接丢弃 return; } MapString, Object doc new HashMap(); doc.put(id, product.getId()); doc.put(title, product.getTitle()); doc.put(price, product.getPrice()); doc.put(createTime, product.getCreateTime()); indexService.saveDoc(product, String.valueOf(product.getId()), doc); } }逻辑说明RabbitListener括号里指定监听的队列名Spring Boot 启动后会自动创建消费者并绑定到该队列。消息处理分三路DELETE直接删 ES 文档不用查 MySQLINSERT和UPDATE共用同一套逻辑回查 MySQL 后写入 ES因为 ES 按主键覆盖写所以在 ES 侧插入和更新没有区别。消费者不捕获异常让它抛出去交给 Spring RabbitMQ 的重试机制处理。这里要注意一个细节SyncMessage必须有无参构造和 getter/setterJackson2JsonMessageConverter反序列化时需要这些基础要素否则会报InvalidDefinitionException。4.2 手动 ack 还是自动 ack同步场景必须手动确认RabbitListener默认是自动 ack即消息一进入消费者方法就确认删除。但同步场景下如果消费者处理失败ES 挂了、MySQL 连接超时消息已经被确认丢掉了ES 就永远少了这条数据。所以要改成手动确认模式spring: rabbitmq: listener: simple: acknowledge-mode: manual配置完之后消费者方法要接收Channel和deliveryTag自己控制确认时机RabbitListener(queues RabbitSyncConfig.SYNC_QUEUE) public void onMessage(SyncMessage message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws IOException { try { // 同步逻辑回查 MySQL、写入 ES channel.basicAck(deliveryTag, false); } catch (Exception e) { channel.basicNack(deliveryTag, false, true); } }参数说明basicAck(tag, false)第二个参数false表示不批量确认只确认当前这条basicNack(tag, false, true)第三个参数true表示重新入队。重试是 RabbitMQ 原生的行为消息回到队列头部继续投递给消费者。要注意死循环问题如果 ES 一直不可用消息会一直入队重试压缩队列吞吐。给重试加一个上限是常见做法用消息头里的x-death统计重试次数超过几次就写入死信队列或者记录日志后 ack 丢弃。4.3 幂等设计重复消费不产生脏数据RabbitMQ 的消息投递语义是「至少一次」配合手动 ack 的重新入队机制消费重复是很正常的。好在 ES 的设计帮了大忙主键相同就覆盖写入所以INSERT和UPDATE消息重复消费ES 里的数据不会重复只是多写一次而已DELETE重复删除也不会报错。真正的幂等风险在别处如果 MySQL 里某条记录被删除消息进入队列此时消费者回查返回null直接 return 不做处理——但 ES 里那条数据其实还在。这种情况需要靠业务约束兜底比如删除操作在业务层同时发两条消息一条DELETE带业务主键消费端直接按主键删 ES 文档。5. 避坑笔记ES 8.3、RabbitMQ、MySQL 三件套里的五个高频问题5.1 ES 8.3 默认开安全认证本地连接一直 401现象Spring Boot 项目启动没问题但一调用 ES 客户端就报401 Unauthorized后台日志显示missing authentication credentials。原因ES 8.0 之后默认开启 xpack 安全认证且默认走 HTTPS不是你以为的裸 HTTP 就能连。本地开发如果不做任何配置ES 只有内置的elastic用户和随机生成的密码客户端没带凭证自然被拒。解决本地测试最简单的方式是 docker 启动时关掉安全认证docker run -d --name es83 -p 9200:9200 -p 9300:9300 \ -e discovery.typesingle-node \ -e xpack.security.enabledfalse \ -e ES_JAVA_OPTS-Xms512m -Xmx512m \ docker.elastic.co/elasticsearch/elasticsearch:8.3.0启动后访问http://localhost:9200会直接返回 JSON客户端配置里也不需要BasicCredentialsProvider。如果必须开启认证就用第 2 章里带用户名密码的配置方式并且把访问地址从https://localhost:9200改成http://localhost:9200因为关掉安全后 HTTPS 也随之失效了。5.2 RabbitMQ 消费者报clean channel shutdown消息堆积不动现象RabbitMQ 管理后台的队列里消息数量持续上涨消费者日志里反复出现channel shutdown或clean channel shutdown但消费者并没有处理消息。原因RabbitListener方法抛出异常后Spring 的监听容器会关闭当前 channel然后重新建立连接。如果异常一直在抛channel 就在「关闭-重建-关闭」里循环消息根本进不了处理方法。最常见的是反序列化失败或者 ES 客户端抛了IOException。解决在消费者方法里对异常类型做区分。反序列化失败属于永久性错误重试多少次都一样直接 ack 丢弃并记录日志ES 暂时不可用、MySQL 连接超时属于暂时性错误basicNack重新入队。再配合spring.rabbitmq.listener.simple.retry.enabledtrue和max-attempts3让 Spring 在内存里完成三次重试三次都失败才走basicNack这样不会频繁触发 channel 重连。5.3 事务内发消息消费者永远查到旧数据现象更新 MySQL 后 ES 里的值没变或者偶尔能更新偶尔不能数据看起来「时灵时不灵」。原因在Transactional方法内部直接调用rabbitTemplate.convertAndSend消息发出时事务可能还没提交消费者立刻回查 MySQL 查到的是事务前的旧值把旧数据写进 ES覆盖了后续的正确更新。解决严格按第 3 章的写法用TransactionSynchronizationManager.registerSynchronization注册afterCommit回调只有事务提交成功才发消息。换一个角度理解MySQL 提交和 MQ 发送这两个动作不能放在一个事务里但必须保证触发顺序是「提交成功 → 发送消息」。这也是为什么后面可以做更进一步的设计——监听 MySQL binlog 的方案从底层就规避了这个问题。5.4 MySQL 的 datetime 字段进入 ES 变成 text范围查询全废现象写入 ES 的createTime字段在 Kibana 里看是2024-01-01 12:00:00但用range查询时报错说字段类型是text不支持排序和范围查询。原因ES 的动态映射对2024-01-01 12:00:00这种字符串的默认推断是text而 text 类型不参与排序和范围聚合。建索引时如果没有显式声明 mapping后面就只能靠数据自己猜猜错了就翻车。解决第 2 章的createIndexIfNotExists里已经做了正确示范——建索引时显式声明.properties(createTime, p - p.date().format(yyyy-MM-dd HH:mm:ss))。如果索引已经建错了没有后悔药只能删掉重建然后全量同步。所以任何 ES 索引第一次创建时必须把 mapping 设计好不要依赖动态映射这是血泪经验。5.5 RabbitMQ 版本和 Erlang 版本不匹配windows / linux 启动失败现象下载了新版 RabbitMQ比如 4.1.x启动时后台日志报 Erlang 版本不支持或者服务起不来管理界面打不开。原因RabbitMQ 对 Erlang 版本有严格的下限要求新版 RabbitMQ 需要对应版本的 Erlang/OTP用系统包管理器随便装的 Erlang 往往是旧版或者路径不对导致 RabbitMQ 诊断脚本找不到运行时。解决安装 RabbitMQ 前先看官方发布的版本兼容矩阵Windows 上直接安装官方配套的 Erlang 版本Linux 上优先用 RabbitMQ 官方提供的二进制包而不是发行版仓库里的版本。装好后用rabbitmqctl status验证 Erlang 版本和节点状态不要急着写 Spring Boot 代码先把http://localhost:15672管理界面打开确认 MQ 本身是健康的再继续。6. 全链路验证技巧一个测试接口跑通 MySQL → MQ → ES 的完整闭环链路级验证不要靠肉眼盯着三个控制台而是写一个组合接口完成「写入 MySQL 查询 ES」的操作闭环。基于前面已经有的ProductService和ProductIndexService补两个接口RestController RequestMapping(/product) public class ProductController { Autowired private ProductService productService; Autowired private ProductIndexService indexService; PostMapping public String save(RequestBody Product product) { productService.save(product); // 业务数据入 MySQL事务提交后自动发 MQ return saved; } GetMapping(/search) public ListObject search(RequestParam String keyword) throws IOException { return indexService.search(product, title, keyword); } }对应的search方法在ProductIndexService里补上public ListObject search(String indexName, String field, String keyword) throws IOException { var response client.search(s - s .index(indexName) .query(q - q.match(m - m.field(field).query(keyword))), Object.class); return response.hits().hits().stream() .map(hit - hit.source()) .toList(); }验证步骤按顺序操作先POST /product插入一条商品数据紧接着GET /product/search?keyword键盘如果 ES 里能搜到刚刚插入的记录说明整条链路已经跑通。如果搜不到按三段排查RabbitMQ 管理后台看queue.mysql.es.sync队列里有没有消息、消费者有没有报错、ES 的product索引里有没有id1的文档。三个位置能快速定位问题出在发送端、消费端还是存储端。这个 demo 值得做做完之后你会有一种「三件套在手里转起来」的感觉。ES 8.3 的版本特性、RabbitMQ 的手动 ack、MySQL 数据到 ES 的字段映射这些单拆开都是文档里能查到的东西但它们之间的串联动起来之后你会理解为什么生产环境要引入消息中间件而不是直接改完 MySQL 就调用 ES 写入——同步逻辑一旦变成异步主流程的响应时间不受影响ES 短暂不可用也不会阻塞业务写入。我第一次跑这个链路时卡在 ES 8 默认安全认证上整整一个晚上后来发现要么关 xpack 要么带凭证就是一行配置的事。事务内发消息的坑是后来才补上的当时线上数据出现了不一致排查到凌晨才发现是事务没提交消息就出去了。希望这些踩过的坑能帮你省下同样的时间也祝你的第一条同步链路一次跑通。本文还有配套的精品资源点击获取
返回列表