首页/新闻资讯/正文详情

Spring Boot + RabbitMQ + Elasticsearch 8.3 数据同步实战:从踩坑到高吞吐

发布时间:2026/9/26 5:13:11 来源:云帆数科 栏目:资讯中心
Spring Boot + RabbitMQ + Elasticsearch 8.3 数据同步实战:从踩坑到高吞吐
简介这是一份面向Java后端开发者的实战Demo聚焦Spring Boot与Elasticsearch 8.3整合并借助RabbitMQ完成MySQL数据向搜索引擎的实时同步适合正在搭建搜索服务或学习消息驱动架构的中高级开发者参考。资源包共78个文件以21个Java源码、21个class编译文件、17个xml配置及2个yml配置文件为主另含证书、构建脚本与说明文档压缩后约90KB结构紧凑便于快速导入IDE运行调试。内容覆盖Spring Data Elasticsearch的Repository操作、RabbitMQ消息监听与发送、数据库变更事件监听、批量写入与错误重试等关键环节并涉及微服务数据流设计与监控思路可帮助读者理解MySQL到Elasticsearch的完整同步链路。目前已有209人学习下载适合作为实时搜索与数据同步场景的入门实践模板。1. 从一次数据对不上的事故说起这套 Spring Boot RabbitMQ Elasticsearch 8.3 的 demo 到底解决什么去年帮一个做 SaaS 工单系统的团队排查问题运营在后台搜一个三天前的工单编号Elasticsearch 里死活搜不到MySQL 里明明躺着这条记录。翻日志才发现他们的同步逻辑是「业务代码里写完 MySQL 顺手调一次 ES 的 REST 接口」某次 ES 集群重启那批写入全丢了业务方还浑然不觉。这个场景几乎每个做搜索的团队都会撞上MySQL 是事实来源Elasticsearch 是查询副本两者之间如果没有一条可靠的消息通道数据一致性就是玄学。这份 demo 要解决的正是这件事——用 Spring Boot 把 MySQL、RabbitMQ、Elasticsearch 8.3 串成一条可复现的同步链路业务侧只写 MySQL变更事件丢进 RabbitMQ消费者负责落到 ES。它适合正在做商品搜索、日志检索、工单查询这类「主库 搜索副本」架构的后端同学也适合想搞明白消息队列在真实项目里怎么落地、而不是停留在「Hello World」的开发者。Elasticsearch 8.x 默认开启安全认证和 7.x 的接入方式差别不小这份 demo 把这块也一并处理了。2. 环境搭建Elasticsearch 8.3、RabbitMQ 与 MySQL 的版本对齐2.1 为什么版本对齐比想象中重要Elasticsearch 8.3 是 8.x 系列里比较稳的一个版本最大的变化是默认启用 HTTPS 和账号密码认证xpack.security.enabled默认就是true。很多人在 Windows 上启动 Elasticsearch 之后发现浏览器访问9200报证书错误或者 Spring Boot 启动时报ElasticsearchRestClientHealthIndicator : Elasticsearch health check failed根因基本都是没处理 8.x 的安全配置。RabbitMQ 这边3.11 之后管理界面的默认账号guest只允许本机登录远程访问必须新建用户并授权 virtual host否则你会遇到「管理界面能打开但 admin 用户创建不了虚拟主机」这种典型问题。MySQL 用 8.0 即可注意caching_sha2_password认证插件和旧版 JDBC 驱动的兼容性。我一般会先把三个中间件的版本固定下来写进docker-compose.yml或者本地安装记录里避免「我本地能跑你那边不行」的扯皮。下面这套版本组合是实测能跑通的组件版本关键配置项Elasticsearch8.3.0xpack.security.enabledtrue需配置证书或关闭 HTTPSRabbitMQ3.11.x新建用户 授权 virtual hostMySQL8.0.x字符集utf8mb4时区Asia/ShanghaiSpring Boot2.7.x对应spring-boot-starter-data-elasticsearchJDK17ES 8.x 客户端要求 JDK 172.2 本地启动 Elasticsearch 8.3 的实操步骤Windows 下启动 Elasticsearch 8.3解压后直接跑bin\elasticsearch.bat会生成一个随机密码打印在控制台同时config\elasticsearch.yml里默认开启了安全。如果你只是本地开发最省事的做法是先把安全关掉等链路跑通再开回来。修改config\elasticsearch.yml# 本地开发临时关闭安全认证生产环境务必开启 xpack.security.enabled: false xpack.security.http.ssl.enabled: false xpack.security.transport.ssl.enabled: false # 单节点模式避免脑裂 discovery.type: single-node改完重启访问http://localhost:9200应该能看到集群信息。这里有个血泪经验如果你之前启动过一次并生成了数据目录关掉安全后可能仍然报认证失败把data目录清掉再启动。Linux 下用docker-compose部署 Elasticsearch 时记得给容器分配足够内存vm.max_map_count要调到 262144否则 ES 起不来。# Linux 宿主机需要调整内核参数 sudo sysctl -w vm.max_map_count262144 # 验证 sysctl vm.max_map_countvm.max_map_count这个参数控制进程可以拥有的内存映射区域数量Elasticsearch 的 Lucene 底层大量使用 mmap默认值 65530 不够用启动时会直接抛max virtual memory areas vm.max_map_count [65530] is too low。这个坑在kubesphere或docker-compose部署 ES 时特别常见很多人以为是镜像问题其实是宿主机内核参数没调。2.3 RabbitMQ 安装与 virtual host 权限配置RabbitMQ 在 Windows 上安装需要先装 Erlang版本要匹配否则rabbitmq-server.bat启动失败。装完之后管理界面默认在15672端口用guest/guest登录。但如果你是在另一台机器上访问guest会被拒绝必须新建用户# 新建用户并设置密码 rabbitmqctl add_user admin YourPassword # 设置为管理员角色 rabbitmqctl set_user_tags admin administrator # 创建虚拟主机 rabbitmqctl add_vhost /demo # 授权 admin 对 /demo 的完整权限 rabbitmqctl set_permissions -p /demo admin .* .* .*这三条权限参数分别对应配置、写、读的正则匹配.*表示全部允许。很多人只建了用户没建 virtual host或者建了 vhost 没授权结果 Spring Boot 连接时报ACCESS_REFUSED - Login was refused using authentication mechanism PLAIN排查半天以为是密码错了。记住RabbitMQ 的权限是「用户 virtual host 权限」三元组缺一不可。3. Spring Boot 整合三件套依赖、配置与核心代码3.1 依赖选型为什么用 spring-boot-starter-data-elasticsearchSpring Boot 2.7 对应的spring-boot-starter-data-elasticsearch底层用的是 Elasticsearch Java API Client 7.17而 ES 8.3 服务端和 7.17 客户端在 HTTP 协议层面是兼容的所以能跑通。但如果你直接用 ES 8.x 的官方 Java Client需要 JDK 17 且 API 写法完全不同。这份 demo 选择前者是为了降低接入成本。pom.xml关键依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-elasticsearch/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.33/version /dependency dependency groupIdcom.baomidou/groupId artifactIdmybatis-plus-boot-starter/artifactId version3.5.3.1/version /dependencyMyBatis-Plus 不是必须的用 JPA 或原生 JDBC 都行这里选它是因为写 CRUD 快demo 里能少写很多样板代码。注意mysql-connector-java8.0.33 支持caching_sha2_password如果你用 5.x 驱动连 MySQL 8会报Unable to load authentication plugin caching_sha2_password。3.2 application.yml 里的连接配置与踩坑点配置文件里最容易翻车的是 Elasticsearch 的地址和 RabbitMQ 的 virtual host。ES 8.3 如果开了安全认证uris要写https并且配置用户名密码如果关了安全写http即可。RabbitMQ 的virtual-host必须和你在rabbitmqctl里建的一致默认是/。spring: elasticsearch: uris: http://localhost:9200 connection-timeout: 5s socket-timeout: 30s rabbitmq: host: localhost port: 5672 username: admin password: YourPassword virtual-host: /demo listener: simple: acknowledge-mode: manual # 手动确认防止消息丢失 prefetch: 10 # 每次预取10条避免消费者过载 datasource: url: jdbc:mysql://localhost:3306/demo?useUnicodetruecharacterEncodingutf8mb4serverTimezoneAsia/Shanghai username: root password: root driver-class-name: com.mysql.cj.jdbc.Driveracknowledge-mode: manual是这份 demo 的关键配置。默认的自动确认模式下消息一投递给消费者就标记为已消费如果消费者处理到一半挂了消息就丢了。手动确认要求你在业务逻辑执行成功后显式调用basicAck失败时basicNack让消息重新入队。prefetch: 10控制的是消费者未确认消息的最大数量设太大容易内存溢出设太小吞吐上不去10 是个比较稳的起步值。3.3 消息生产者MySQL 写入后如何可靠地发消息核心思路是「先写库再发消息」但这两步不在一个事务里所以要么用本地消息表要么接受极小概率的不一致。demo 里用的是简化版在 Service 层写完 MySQL 后直接发 RabbitMQ如果发送失败就记录日志人工补偿。生产环境建议上本地消息表或事务消息。Service public class ProductService { Autowired private ProductMapper productMapper; Autowired private RabbitTemplate rabbitTemplate; // 写库后发送消息到交换机 public void saveProduct(Product product) { // 1. 先落 MySQL productMapper.insert(product); // 2. 构造消息体带上操作类型 ProductMessage msg new ProductMessage(); msg.setId(product.getId()); msg.setType(INSERT); msg.setTimestamp(System.currentTimeMillis()); // 3. 发送到指定交换机和路由键 rabbitTemplate.convertAndSend( product.exchange, // 交换机名称 product.insert, // 路由键 msg ); } }convertAndSend的第一个参数是交换机第二个是路由键第三个是消息体。消息体默认用 JDK 序列化跨语言场景建议换成 JSON配置一个Jackson2JsonMessageConverter即可。这里有个细节如果 MySQL 插入成功但convertAndSend抛异常数据就不一致了。常见做法是把消息先写到一张message_log表再用定时任务扫描重发这就是本地消息表模式。3.4 消息消费者手动 ACK 与 ES 批量写入消费者这边要做三件事接收消息、写 ES、手动 ACK。写 ES 用ElasticsearchRestTemplate或者RestHighLevelClientdemo 里用前者更简洁。Component public class ProductConsumer { Autowired private ElasticsearchRestTemplate esTemplate; Autowired private ProductMapper productMapper; RabbitListener(queues product.queue) public void handleMessage(ProductMessage msg, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws IOException { try { // 1. 根据操作类型处理 if (INSERT.equals(msg.getType()) || UPDATE.equals(msg.getType())) { Product product productMapper.selectById(msg.getId()); if (product ! null) { // 2. 写入 ESid 用 MySQL 主键保证幂等 IndexQuery query new IndexQueryBuilder() .withId(String.valueOf(product.getId())) .withObject(product) .build(); esTemplate.index(query, IndexCoordinates.of(product_index)); } } else if (DELETE.equals(msg.getType())) { esTemplate.delete(String.valueOf(msg.getId()), Product.class); } // 3. 处理成功手动 ACK channel.basicAck(tag, false); } catch (Exception e) { // 处理失败拒绝消息并重新入队或进死信队列 channel.basicNack(tag, false, false); } } }basicAck的第二个参数multiple设为false表示只确认当前这条。basicNack的第三个参数requeue设为false表示不重新入队直接进死信队列避免无限重试打爆消费者。这里用 MySQL 主键作为 ES 文档 id天然幂等同一条消息重复消费也不会产生脏数据。如果你用 ES 自动生成 id重复消费就会插入两条相同记录这是新手最容易踩的坑之一。4. 避坑与排查同步链路上最常见的五个翻车现场4.1 现象ES 里搜不到刚写入的数据原因通常是消息没发出去或者消费者没消费。先看 RabbitMQ 管理界面product.queue的Ready和Unacked数量。如果Ready一直涨说明消费者没启动或者消费能力不足如果Unacked一直涨说明消费者卡住了大概率是 ES 写入超时。解决检查消费者日志有没有ElasticsearchException确认 ES 集群健康状态GET /_cluster/health是不是yellow或red。4.2 现象消费者无限重试同一条消息原因是basicNack的requeue设成了true而消息本身有格式问题或依赖数据不存在每次重试都失败。解决把requeue设为false配置死信交换机DLX让失败消息进死信队列人工排查后再决定是否重投。死信队列的配置在application.yml里声明队列时绑定x-dead-letter-exchange参数即可。4.3 现象Spring Boot 启动报 Elasticsearch health check failed这是热词里出现频率很高的问题。原因通常是 ES 8.3 开了安全认证但 Spring Boot 没配用户名密码或者用了https但证书没信任。解决本地开发先关掉xpack.security.enabled生产环境则配置spring.elasticsearch.username和password并把 CA 证书导入 JDK 的cacerts。如果只是不想让健康检查影响启动可以配置management.health.elasticsearch.enabledfalse临时绕过。4.4 现象RabbitMQ 管理界面能打开但 admin 用户创建不了虚拟主机这是权限问题。admin用户虽然角色是administrator但如果没有对目标 vhost 的权限仍然无法操作。解决用rabbitmqctl set_permissions -p /demo admin .* .* .*授权注意-p后面跟的是 vhost 名称。另外administrator角色只代表能登录管理界面不代表对所有 vhost 有读写权限这两个概念经常被混淆。4.5 现象MySQL 数据更新了但 ES 里还是旧值原因是更新操作没有发消息或者消息里的type字段没匹配上消费者的分支。解决检查 Service 层的更新方法有没有调用convertAndSend以及消费者里if判断的字符串是否和生产者一致。建议把操作类型定义成枚举避免手写字符串拼错。另外如果用了 MyBatis-Plus 的updateById注意它返回的是影响行数不是更新后的实体消费者需要重新查一次库。5. 进阶技巧用批量消费和索引模板把同步吞吐拉上去单条消费在数据量小的时候没问题一旦 MySQL 每秒几百条变更逐条写 ES 会成为瓶颈。我一般会把消费者改成批量模式配合prefetch和手动 ACK一次拉一批消息攒够 50 条或超过 1 秒就批量写 ES。ElasticsearchRestTemplate的bulkIndex方法支持批量操作比逐条index快一个数量级。// 批量消费示例攒够50条或超时1秒触发批量写入 RabbitListener(queues product.queue, containerFactory batchFactory) public void handleBatch(ListProductMessage messages, Channel channel) throws IOException { ListIndexQuery queries new ArrayList(); long maxTag 0; for (ProductMessage msg : messages) { Product product productMapper.selectById(msg.getId()); if (product ! null) { queries.add(new IndexQueryBuilder() .withId(String.valueOf(product.getId())) .withObject(product) .build()); } } if (!queries.isEmpty()) { esTemplate.bulkIndex(queries, IndexCoordinates.of(product_index)); } // 批量确认最后一条的 tagmultipletrue 表示确认之前所有 channel.basicAck(maxTag, true); }批量消费的containerFactory需要单独配置设置consumerBatchEnabledtrue和batchSize。basicAck的multiple设为true时会确认 delivery tag 小于等于当前 tag 的所有消息所以批量场景下只需要确认最后一条。这里有个后悔药级别的提醒批量写入 ES 时如果部分文档失败bulkIndex不会抛异常而是返回一个包含错误信息的响应你必须手动检查BulkResponse.hasFailures()否则会静默丢数据。另一个进阶点是索引模板。ES 8.3 支持 composable index templates可以在写入前定义好 mapping避免动态映射把数字字段识别成 text。比如商品价格字段如果不预设 mappingES 可能把99.9映射成text导致范围查询失效。用Document(indexName product_index)注解配合Field(type FieldType.Double)可以解决但更稳妥的做法是在 ES 侧用模板统一管理。PUT _index_template/product_template { index_patterns: [product_*], template: { mappings: { properties: { id: {type: long}, name: {type: text, analyzer: ik_max_word}, price: {type: double}, createTime: {type: date, format: yyyy-MM-dd HH:mm:ss} } } } }ik_max_word是中文分词插件需要单独安装不装的话中文搜索会按单字切分效果很差。索引模板的好处是新建索引时自动套用 mapping不用每次手动建。从那以后我每次接搜索需求都强制先确认三件事ES 版本和安全配置、RabbitMQ 的 vhost 权限、以及索引 mapping 有没有预设。这三步走完后面基本不会出大问题。希望帮到你。本文还有配套的精品资源点击获取

相关推荐

Mac上用Docker搭建MySQL:从环境配置到数据持久化完整指南
Mac上用Docker搭建MySQL:从环境配置到数据持久化完整指南

作为一个长期在 Mac 上写后端代码的人,我换过好几台 MacBook,每次重装系统最烦的事情之一就是数据库环境。以前我习惯用brew install mysql直接装,后来发现项目多了之后,版本冲突、依赖残留、卸载不干净这些问题轮番折磨人。直到我… · 2026/9/26 5:13:11

Objective-C短信接口对接实战:从SDK选型到上架审核
Objective-C短信接口对接实战:从SDK选型到上架审核

做iOS开发的老哥,十有八九都接过这种需求:App注册页面要加手机号验证码登录,或者找回密码必须走短信验证,再不然就是运营那边提了个需求,要给用户发通知短信。短信接口这东西技术门槛不算高,但坑是真的多&a… · 2026/9/26 5:13:11

大功率超声波口罩点焊机详解:20kHz与15kHz选型调试与维护指南
大功率超声波口罩点焊机详解:20kHz与15kHz选型调试与维护指南

1. 大功率超声波口罩点焊机:先认清这套系统的构成与工作原理前几年口罩设备行情最火爆的时候,我前前后后经手过几十台口罩点焊机,其中20kHz和15kHz这两类大功率超声波设备占了一大半。这类机器结构看着不复杂,但不少用户买回去只在… · 2026/9/26 5:13:05

ThinkPHP+Laravel+Vue二手车销售平台开发实战
ThinkPHP+Laravel+Vue二手车销售平台开发实战

做二手汽车销售平台,一开始摆在面前的两条路就挺有意思。项目标题里同时挂了ThinkPHP和Laravel,很多同行看到第一反应是“这俩框架选一个不就完了吗”。实际做下来你会发现,真正落地的项目里,这个选择题背后牵扯的是团队技术栈、服… · 2026/9/26 7:56:47

UE5内置建模工具链:Modeling Mode与Geometry Script实战指南
UE5内置建模工具链:Modeling Mode与Geometry Script实战指南

1. 从“37”说起:为什么 UE5 的建模工具链值得单独拎出来聊 如果你最近在 UE5 里折腾过场景搭建,大概率会遇到一个尴尬的瞬间:美术给的模型还没到位,但你想先摆个白模看看比例;或者从商城买来的资产面数爆炸&#xff0… · 2026/9/26 7:56:47

无畏契约Vanguard启动报错全解析:从服务到驱动的排查与修复指南
无畏契约Vanguard启动报错全解析:从服务到驱动的排查与修复指南

1. 先搞清楚Vanguard到底在干什么很多人一看到无畏契约启动报错,第一反应就是“游戏坏了”,然后开始重装游戏、重装系统,折腾一整天问题还在。实际上,无畏契约的启动链路比大多数游戏复杂得多,它不是一个单纯的游戏客户… · 2026/9/26 7:56:35

iOS国密改造实战:OpenSSL集成SM2/SM4与避坑指南
iOS国密改造实战:OpenSSL集成SM2/SM4与避坑指南

简介:面向iOS平台国密算法开发者的实践参考,内容围绕SM2加密在iOS侧的落地展开,基于GmSSL改造整理,弥补了网上iOS端缺少可直接参考国密示例的空白。作者在C语言基础较弱、现有实现代码杂乱且缺少注释的条件下反复踩坑,… · 2026/9/26 7:56:35

手写SQL解析器:词法分析、AST与生产级选型实践
手写SQL解析器:词法分析、AST与生产级选型实践

简介:基于Flex与Bison这两款开源编译器工具构建的SQL解析器完整工程,面向数据库内核研发和编译器技术学习者,提供从SQL语句输入到词法切分、语法检查、抽象语法树构建再到中间表示输出的完整实现参考。压缩包共包含11个文件,以四个… · 2026/9/26 7:56:29

金融技术服务项目启动前提与内容规范
金融技术服务项目启动前提与内容规范

我无法根据当前输入生成符合要求的博文。原因如下:项目标题为"financial-services",这是一个高度泛化的行业术语,本身不构成具体可操作、可拆解的项目或技术主题;项目正文为空,未提供任何实质性描述、功能定… · 2026/9/26 7:56:29

数据库课后习题答案别硬背:当测试用例集刷,效率翻倍
数据库课后习题答案别硬背:当测试用例集刷,效率翻倍

简介:万常选版《数据库原理与设计》课后习题答案资源,覆盖第2至6章及第9章,适合正在学习关系模型、数据库建模、关系数据理论与模式求精的本科生、自学者作为复习与自测材料。压缩包共7个文件,含3个doc参考答案、2个sql示例脚本、… · 2026/9/26 0:00:21

OpenClaw 替代品?Hermes Agent 踩坑实录:macOS 飞书接入 TaoToken 配置
OpenClaw 替代品?Hermes Agent 踩坑实录:macOS 飞书接入 TaoToken 配置

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/26 0:00:40

向下兼容与向上兼容:接口设计中的兼容性策略与工程实践
向下兼容与向上兼容:接口设计中的兼容性策略与工程实践

一次版本升级事故,是很多团队绕不过去的坎。线上环境里,服务端明明已经上线了新版接口,老的移动端还在照着旧文档传参数。请求一到网关,校验直接拒绝,用户操作失败,客服群炸了锅,开发群里开始互… · 2026/9/26 0:00:46

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码