
简介本资源是一份Spring Boot整合Elasticsearch 8.3并通过RabbitMQ同步MySQL数据的实战Demo面向具备一定Java与Spring基础、希望掌握实时数据同步与搜索服务集成的后端开发者。内容涵盖Spring Data Elasticsearch的Repository操作、RabbitMQ消息监听与发送、MySQL变更事件监听以及MySQL→RabbitMQ→Spring Boot→Elasticsearch的完整数据流链路并涉及批量写入、错误重试、监控与权限控制等优化思路。压缩包共78个文件以21个java源码、21个class编译文件、17个xml配置为主另含yml配置、证书与Maven包装脚本等整体约90KB结构紧凑便于快速导入运行。目前已有209人学习下载适合作为构建实时搜索与数据同步系统的参考模板帮助读者理解消息队列解耦、索引同步与性能调优的落地方式。1. 从一次数据对不上的事故说起这套 Demo 到底解决什么问题凌晨两点运营在群里甩了张截图后台明明显示订单已支付搜索页却死活搜不到。翻日志发现 MySQL 里数据早就落库了Elasticsearch 那边还是旧状态。这种「数据库和搜索引擎各说各话」的场面做过搜索业务的同学应该都不陌生。根子在于双写代码里先写 MySQL 再写 ES中间任何一步抖动——网络超时、ES 集群短暂不可用、事务回滚——两边就永久不一致了。这份 Demo 给的思路很直接MySQL 只负责落库Elasticsearch 的同步交给 RabbitMQ 异步解耦。业务代码写完 MySQL 发一条消息到 MQ消费者拿到消息再去写 ES。这样即使 ES 挂了消息还在队列里堆着等它恢复再消费数据最终一致。技术栈是 Spring Boot RabbitMQ Elasticsearch 8.3适合正在做商品搜索、订单检索、日志聚合这类场景又不想上 Canal、Logstash 全套重装备的团队。下面我按「环境怎么搭 → 消息怎么发怎么收 → ES 8.3 的 API 怎么调 → 坑在哪」的顺序拆一遍能直接抄作业。2. 环境搭建与依赖版本对齐别让版本号成为第一道坎2.1 三个组件的版本选择逻辑Elasticsearch 8.x 最大的变化是默认开启安全认证9200端口不再是裸奔状态。Demo 里用的是 8.3这个版本已经稳定支持 Java API Client就是那个co.elastic.clients包旧的RestHighLevelClient在 7.15 就被标记废弃了。所以选型上有个硬约束Spring Boot 版本不能太低否则自带的依赖管理会把 ES 客户端锁在老版本。我一般这么配Spring Boot 用 3.0.x 或 3.1.x它默认管理的 Elasticsearch Java Client 版本和 8.3 能对上RabbitMQ 用 3.11 以上镜像队列和 quorum 队列都成熟了。JDK 至少 17因为 Spring Boot 3 强制要求。三个组件的版本关系可以用一句话记Spring Boot 管依赖版本ES 管数据存储RabbitMQ 管消息投递三者版本各自独立但客户端要兼容。2.2 pom.xml 关键依赖与排除项dependencies !-- Web 基础提供 REST 接口 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency !-- RabbitMQ 消息队列 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency !-- Elasticsearch Java API Client8.x 新客户端 -- dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version8.3.0/version /dependency !-- 必须显式引入否则启动报 ClassNotFoundException -- dependency groupIdcom.fasterxml.jackson.core/groupId artifactIdjackson-databind/artifactId /dependency !-- MySQL 驱动 -- dependency groupIdcom.mysql/groupId artifactIdmysql-connector-j/artifactId scoperuntime/scope /dependency !-- MyBatis-Plus 或 JPA 二选一Demo 用 MyBatis-Plus -- dependency groupIdcom.baomidou/groupId artifactIdmybatis-plus-boot-starter/artifactId version3.5.3.1/version /dependency /dependencies这里有个容易翻车的点elasticsearch-java和jackson-databind的版本要匹配。8.3 的客户端内部用 Jackson 做序列化如果项目里其他依赖把 Jackson 降到了 2.12 以下启动时就会报NoSuchMethodError。解决办法是在dependencyManagement里统一锁 Jackson 版本到 2.14 以上。2.3 application.yml 的连接配置spring: datasource: url: jdbc:mysql://localhost:3306/demo_db?useSSLfalseserverTimezoneAsia/Shanghai username: root password: your_password driver-class-name: com.mysql.cj.jdbc.Driver rabbitmq: host: localhost port: 5672 username: guest password: guest # 开启手动 ack防止消息丢失 listener: simple: acknowledge-mode: manual prefetch: 10 elasticsearch: host: localhost port: 9200 username: elastic password: your_es_password # 8.x 默认 https本地测试可关掉 scheme: http参数说明acknowledge-mode: manual是关键自动 ack 在消费者抛异常时消息会直接丢手动 ack 才能配合重试。prefetch: 10控制每个消费者预取的消息数太大容易导致消息堆积在单个消费者内存里太小又影响吞吐10 是个折中值。ES 的scheme如果集群开了 TLS 就改成https同时要配证书路径本地开发图省事一般关掉。提示ES 8.3 首次启动会自动生成elastic用户密码在日志里搜Password for the elastic user能找到。忘了的话用bin/elasticsearch-reset-password -u elastic重置。3. RabbitMQ 消息投递与消费同步链路的核心3.1 交换机、队列与路由键的设计同步链路的消息模型不复杂但设计错了后面全是坑。Demo 里用的方案是一个 topic 交换机 一个队列 一个路由键按业务表名区分路由键比如sync.product、sync.order。这样以后要加新表只要新增路由键和消费者就行不用动交换机。Configuration public class RabbitConfig { // 业务同步交换机topic 类型支持通配符路由 public static final String SYNC_EXCHANGE sync.exchange; public static final String SYNC_QUEUE sync.queue; public static final String SYNC_ROUTING_KEY sync.#; Bean public TopicExchange syncExchange() { // durabletrue 表示交换机持久化重启不丢 return new TopicExchange(SYNC_EXCHANGE, true, false); } Bean public Queue syncQueue() { // 第二个参数 durabletrue队列持久化 return new Queue(SYNC_QUEUE, true); } Bean public Binding syncBinding() { return BindingBuilder.bind(syncQueue()) .to(syncExchange()) .with(SYNC_ROUTING_KEY); } }逻辑说明TopicExchange的sync.#能匹配sync.product、sync.order.insert等所有以sync.开头的路由键。durabletrue保证 RabbitMQ 重启后交换机和队列还在但注意消息本身要持久化还得在发送时设deliveryMode2这两件事经常被混为一谈。3.2 生产者写完 MySQL 再发消息生产者的核心原则是先落库再发消息顺序反了会出现「消息发了但事务回滚」的脏数据。Demo 里用Transactional包住数据库操作消息发送放在事务提交后。Service public class ProductService { Autowired private ProductMapper productMapper; Autowired private RabbitTemplate rabbitTemplate; Transactional(rollbackFor Exception.class) public void saveProduct(Product product) { // 1. 先写 MySQL productMapper.insert(product); // 2. 事务提交后再发消息用 TransactionSynchronization 保证时序 TransactionSynchronizationManager.registerSynchronization( new TransactionSynchronization() { Override public void afterCommit() { // 消息体带上操作类型和主键消费者据此决定 insert/update/delete SyncMessage msg new SyncMessage(INSERT, product.getId(), product); rabbitTemplate.convertAndSend( RabbitConfig.SYNC_EXCHANGE, sync.product, msg, message - { // 消息持久化防止 MQ 重启丢消息 message.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT); return message; } ); } } ); } }参数说明afterCommit()是事务提交后的回调确保消息发出时数据已经落库。SyncMessage里带operation字段是为了让消费者知道该做插入还是删除只传数据体的话删除操作没法表达。MessageDeliveryMode.PERSISTENT对应投递模式 2配合队列持久化才能真正做到消息不丢。3.3 消费者手动 ack 与幂等处理消费者这边最容易出两个问题一是 ack 时机不对导致消息重复消费二是没做幂等导致 ES 里出现重复文档。Component public class SyncConsumer { Autowired private ElasticsearchService esService; RabbitListener(queues RabbitConfig.SYNC_QUEUE) public void handleSync(SyncMessage msg, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException { try { // 根据操作类型分发 switch (msg.getOperation()) { case INSERT: case UPDATE: esService.upsertDocument(msg.getData()); break; case DELETE: esService.deleteDocument(msg.getId()); break; } // 业务处理成功才 ack channel.basicAck(tag, false); } catch (Exception e) { // 处理失败requeuefalse 进死信队列避免无限重试 channel.basicNack(tag, false, false); log.error(同步 ES 失败, id{}, msg.getId(), e); } } }逻辑说明basicAck放在业务逻辑之后保证「处理成功才确认」。basicNack的第三个参数requeuefalse表示不重新入队消息会进死信队列如果配了的话否则失败消息会无限循环重试把队列堵死。幂等靠 ES 的_id实现——用 MySQL 主键作为 ES 文档 ID重复写入就是覆盖不会产生重复文档。注意prefetch和手动 ack 配合时如果消费者处理慢未 ack 的消息会一直占着 prefetch 配额。生产环境建议把 prefetch 设成 1~10 之间别贪大。4. Elasticsearch 8.3 客户端操作新 API 的写法与参数4.1 Java API Client 的初始化8.x 的客户端初始化和 7.x 完全不同RestHighLevelClient那套已经不能用了。新客户端基于ElasticsearchTransport构建配置项通过RestClient传入。Configuration public class EsConfig { Value(${elasticsearch.host}) private String host; Value(${elasticsearch.port}) private int port; Value(${elasticsearch.username}) private String username; Value(${elasticsearch.password}) private String password; Bean public ElasticsearchClient elasticsearchClient() { // 底层 REST 客户端负责 HTTP 连接 RestClient restClient RestClient.builder( new HttpHost(host, port, http)) .setRequestConfigCallback(builder - builder.setConnectTimeout(5000) // 连接超时 5s .setSocketTimeout(60000)) // 读超时 60s .build(); // 传输层负责 JSON 序列化 ElasticsearchTransport transport new RestClientTransport( restClient, new JacksonJsonpMapper()); return new ElasticsearchClient(transport); } }参数说明connectTimeout是建立 TCP 连接的超时socketTimeout是等待响应的超时。ES 做聚合查询时响应可能很慢socketTimeout设太小会频繁报SocketTimeoutException60 秒是个安全值。JacksonJsonpMapper负责 Java 对象和 JSON 的互转如果实体类里有LocalDateTime字段需要额外注册JavaTimeModule否则序列化会报错。4.2 索引创建与文档 upsertDemo 里用upsert而不是index因为同步场景下同一条数据可能被多次投递upsert能保证「存在则更新不存在则插入」。Service public class ElasticsearchService { Autowired private ElasticsearchClient client; private static final String INDEX_NAME product_index; // 创建索引指定分片和映射 public void createIndexIfNotExists() throws IOException { boolean exists client.indices().exists(e - e.index(INDEX_NAME)).value(); if (!exists) { client.indices().create(c - c .index(INDEX_NAME) .settings(s - s .numberOfShards(3) // 主分片数建索引后不可改 .numberOfReplicas(1)) // 副本数可动态调整 .mappings(m - m .properties(id, p - p.long_(l - l)) .properties(name, p - p.text(t - t .analyzer(ik_max_word) // 中文分词 .fields(keyword, f - f.keyword(k - k)))) .properties(price, p - p.double_(d - d)) .properties(createTime, p - p.date(d - d .format(yyyy-MM-dd HH:mm:ss)))) ); } } // upsert文档存在则更新不存在则插入 public void upsertDocument(Product product) throws IOException { client.update(u - u .index(INDEX_NAME) .id(String.valueOf(product.getId())) // 用 MySQL 主键做文档 ID .doc(product) .docAsUpsert(true), // 关键不存在时按 doc 插入 Product.class); } public void deleteDocument(Long id) throws IOException { client.delete(d - d.index(INDEX_NAME).id(String.valueOf(id))); } }逻辑说明numberOfShards建索引后不能改所以一开始就要估算数据量。单分片建议 30~50GB数据量不大就 1~3 个分片。docAsUpsert(true)是 upsert 的核心参数少了它文档不存在时会报document_missing_exception。中文分词用了ik_max_word需要提前装 IK 分词插件没装的话用standard也能跑只是中文搜索效果差。4.3 搜索查询的构建public ListProduct search(String keyword, int page, int size) throws IOException { SearchResponseProduct response client.search(s - s .index(INDEX_NAME) .from(page * size) // 分页起始位置 .size(size) // 每页条数 .query(q - q .multiMatch(m - m .fields(name, description) // 多字段匹配 .query(keyword) .type(TextQueryType.BestFields))), Product.class); return response.hits().hits().stream() .map(Hit::source) .collect(Collectors.toList()); }参数说明from size的分页方式在深分页时性能很差from超过 10000 会直接报错。生产环境深分页要用search_after或scroll。multiMatch的BestFields类型表示哪个字段匹配度最高就用哪个字段的得分适合商品名和描述分开存的场景。5. 避坑与排查那些让我加班到凌晨的问题5.1 消息重复消费导致 ES 文档重复现象ES 里同一条商品出现多条文档_id不同但内容一样。原因消费者处理完业务逻辑后、ack 之前进程挂了RabbitMQ 认为消息没被确认重新投递。如果 ES 文档 ID 用的是自动生成的就会产生重复。解决文档 ID 强制用 MySQL 主键upsert天然幂等。另外消费者里加一层 Redis 去重也行但主键做 ID 是最省事的方案。5.2 ES 8.x 连接报 SSL 握手失败现象启动时报SSLException: Unrecognized SSL message, plaintext connection?。原因ES 8.x 默认开启 TLS客户端却用http去连9200。解决要么在elasticsearch.yml里把xpack.security.http.ssl.enabled设为false仅限本地开发要么客户端配https并导入 CA 证书。生产环境别关 TLS。5.3 消费者无限重试把队列堵死现象队列里消息数只增不减日志里同一条错误反复刷。原因basicNack的requeue参数设成了true失败消息重新入队又被消费死循环。解决requeuefalse配合死信队列失败消息进 DLX人工排查后再决定是否重投。同时给队列设x-message-ttl和x-max-length防止无限堆积。5.4 事务未提交就发消息导致数据不一致现象ES 里有数据MySQL 里查不到。原因消息发送写在Transactional方法体内事务还没提交消息就发出去了消费者去查 MySQL 查不到。解决用TransactionSynchronizationManager.registerSynchronization把发送逻辑放到afterCommit回调里确保事务提交后再发。5.5 批量同步时 ES 写入性能骤降现象单条同步没问题批量导入几万条时 ES 响应越来越慢。原因每条消息单独调一次updateAPI网络往返和 Lucene 段合并开销叠加。解决消费者里攒批用BulkRequest批量提交每 500~1000 条刷一次。同时把索引的refresh_interval临时调成-1导入完再调回来。6. 进阶技巧用 Bulk 批处理把同步吞吐拉起来单条 upsert 在数据量上来之后就是瓶颈。我实测过单条同步 QPS 大概在 200 左右换成 Bulk 之后能到 2000 以上差距主要在网络往返和段合并次数上。下面这个消费者改造方案是我现在项目里在用的核心思路是攒批 定时刷 手动 ack。Component public class BatchSyncConsumer { Autowired private ElasticsearchClient client; // 线程安全的缓冲队列 private final BlockingQueueSyncMessage buffer new LinkedBlockingQueue(5000); private static final int BATCH_SIZE 500; private static final long FLUSH_INTERVAL_MS 2000; // 定时刷盘防止消息量少时一直不提交 Scheduled(fixedDelay FLUSH_INTERVAL_MS) public void flushByTime() { if (!buffer.isEmpty()) { doBulkFlush(); } } RabbitListener(queues RabbitConfig.SYNC_QUEUE) public void handle(SyncMessage msg, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException { buffer.offer(msg); if (buffer.size() BATCH_SIZE) { doBulkFlush(); } // 注意这里 ack 的是「已入缓冲」不是「已写 ES」 channel.basicAck(tag, false); } private synchronized void doBulkFlush() { ListSyncMessage batch new ArrayList(); buffer.drainTo(batch, BATCH_SIZE); if (batch.isEmpty()) return; try { BulkRequest.Builder br new BulkRequest.Builder(); for (SyncMessage msg : batch) { if (DELETE.equals(msg.getOperation())) { br.operations(op - op.delete(d - d .index(product_index) .id(String.valueOf(msg.getId())))); } else { br.operations(op - op.update(u - u .index(product_index) .id(String.valueOf(msg.getId())) .doc(msg.getData()) .docAsUpsert(true))); } } BulkResponse response client.bulk(br.build()); if (response.errors()) { // 逐条检查失败项记录后人工处理 response.items().stream() .filter(item - item.error() ! null) .forEach(item - log.error(Bulk 失败: {}, item.error().reason())); } } catch (Exception e) { log.error(Bulk 提交异常, e); } } }这段代码有几个设计取舍值得说清楚。第一ack 时机提前到了入缓冲之后这意味着如果进程在刷盘前挂了缓冲里的消息会丢。要解决这个问题可以把 ack 放到doBulkFlush成功之后但那样就得自己维护 deliveryTag 和消息的对应关系复杂度上来了。我一般根据业务容忍度选搜索场景丢几条能接受就用前者订单场景必须用后者。第二Scheduled的定时刷盘是兜底防止低峰期消息在缓冲里躺太久。第三synchronized保证定时任务和消费者线程不会同时刷盘避免重复提交。验证同步是否生效我习惯用两步先在 MySQL 里改一条数据等两秒然后直接查 ES 的_doc接口看文档版本号有没有变。命令是curl -u elastic:密码 http://localhost:9200/product_index/_doc/1返回的_version递增就说明同步成功。如果版本号没动先看 RabbitMQ 管理页面的队列有没有堆积再看消费者日志有没有异常最后才怀疑 ES 本身。从那以后我每次接同步链路都强制先跑一遍「改 MySQL → 查 ES 版本号」这个最小验证确认链路通了再往上堆业务逻辑。希望帮到你。本文还有配套的精品资源点击获取