FEATURED · 精选文章

MySQL数据同步到Elasticsearch:基于RabbitMQ的异步解耦方案

发布时间 / 2026/9/8 6:45:54
来源 / 创域科博编辑部
栏目 / 资讯中心
MySQL数据同步到Elasticsearch:基于RabbitMQ的异步解耦方案 简介一份基于Spring Boot与Elasticsearch 8.3的集成Demo面向需要实现MySQL数据与搜索引擎实时同步的Java后端开发者也可作为微服务架构中异步数据流转的入门参考。资源演示了通过RabbitMQ消息队列异步处理数据库增删改事件并同步写入Elasticsearch的完整流程涵盖Spring Data Elasticsearch依赖与连接配置、Repository接口定义、RabbitMQ消息监听器、JDBC数据源接入以及使用TransactionalEventListener监听数据库变更等关键环节。压缩包仅90KB共78个文件以21个Java源码、21个编译后class、17个XML配置及2个YAML配置文件为主其中XML多为Maven依赖与MyBatis映射YAML为Spring Boot与连接配置另含Maven wrapper脚本、Git忽略文件和HELP说明文档目录结构清晰便于按模块阅读调试。已有208人学习浏览。借助该Demo可快速搭建MySQL→RabbitMQ→Spring Boot→Elasticsearch的数据同步原型学习批量写入、消息重试与错误处理等优化手法测试方面还可参考其单元测试与监控设计思路对构建电商搜索、日志分析等实时检索场景具有直接借鉴价值。1. 项目概述1.1 这个 Demo 到底解决了什么问题先聊一个最常见的业务场景你做电商、内容社区或者企业级应用核心数据存在 MySQL 里商品表、文章表、订单表都在。有一天产品经理说“搜索太慢了要支持分词、支持高亮、支持复杂的过滤查询”你发现 MySQL 的LIKE %关键词%在数据量大了之后慢性死亡索引失效、全文索引又难用。这时候 Elasticsearch 就该登场了。但问题紧接着就来了——ES 不能直接替代 MySQL它更像一个加速器、检索层。业务数据的主存储和事务管理还在 MySQLES 只是把需要被检索的字段同步过去一份。数据怎么同步最粗暴的办法是双写业务代码里更新完 MySQL 再调一次 ES API 写入。听起来简单但一旦 ES 集群抖动或者写入失败业务流程会被拖垮而且这样做让业务代码变得很脏。这个 Demo 的核心思路就是引入 RabbitMQ 作为中间层。MySQL 的数据发生增删改之后把变更事件丢进消息队列由独立的消费端去异步更新 ES。这样做的好处有三个一是业务写入链路上不直接依赖 ES 的可用性生产环境里 ES 集群故障不会拖垮主业务流程二是削峰大促时大量更新请求不至于瞬间打爆 ES 的写入能力三是同一份变更消息可以被多个消费者重复使用比如说既更新搜索索引也更新推荐系统的数据做数据分发非常方便。1.2 技术选型和版本说明这套 Demo 用到的技术栈如下Spring Boot 2.7.x为什么不选 3.x因为 3.x 要求 JDK 17 起步很多公司的存量项目还在 8 或者 11本 Demo 面向大众用了更稳妥的组合。Elasticsearch 8.3.x8.x 版本变化很大内置了 Java 客户端的新 API废除了 RestHighLevelClient很多老教程已经失效这也是写这个 Demo 的一个重要原因。RabbitMQ 3.10.x Erlang 24/25这是比较成熟的搭配。MySQL 8.0.x存储业务原始数据。JDK 8 或 11实测都能跑。我把整套 Demo 源码全部走通了一遍包括环境搭建、索引创建、消息发送、消费同步。接下来我会把每一步的细节、坑点以及最终的代码骨架都拆开来讲保证你照着敲一遍就能在本地把数据从 MySQL 流转到 ES然后在 Elasticsearch 里搜到它。注意ES 8.3 对 JDK 有硬性要求至少 JDK 17 才能启动 Elasticsearch 服务本身但 Spring Boot 项目可以用 JDK 8 开发两者并不冲突很多人在这里踩坑后面环境准备部分我会单独展开。2. 环境准备与组件安装2.1 Elasticsearch 8.3 Windows 安装的特殊之处网上 Elasticsearch 安装教程很多但大部分讲的是 7.x 版本。8.x 有几个显著变化初次接触的人很容易卡住。第一JDK 版本。ES 8.3 内置了 JDK 17如果你直接下载压缩包解压它会使用自带 JDK 启动不需要你额外装 JDK 17。但如果你设置过JAVA_HOME环境变量指向 JDK 8启动时就会报错需要在config/jvm.options里加一行-Djava.io.tmpdir${ES_TMPDIR}或者直接修改启动脚本。我的建议是如果你开发环境里同时有多个 JDK 版本最好在启动 ES 前临时把JAVA_HOME指向 JDK 17 的路径或者完全不设置JAVA_HOME让它用内置版本。第二安全认证默认开启。8.x 首次启动会生成一个elastic超级用户的密码显示在控制台上还会生成证书文件。如果你忘记保存初始密码后面访问 ES 会非常麻烦。进入config/elasticsearch.yml可以看到xpack.security.enabled: true对于本地开发如果你觉得麻烦可以把认证关掉设置成false。但这会失去 8.x 的一个新增安全特性而且生产环境千万不要这样干。我建议搞懂原理之后在 Demo 里保留认证代码里配置账号密码访问这样贴近真实项目。第三端口变化。7.x 时代你可能习惯用 9200 做 HTTP、9300 做集群内部通信。8.3 虽然还保留这些端口但新装集群的节点间通信默认使用 9300 端口且强制要求 TLS 证书这点和 7.x 反而一样。本地单节点开发不用纠结记住访问地址是https://localhost:9200或http://localhost:9200看你http.type配置。2.2 RabbitMQ 安装和必备插件RabbitMQ 是 Erlang 写的所以安装顺序是先装 Erlang再装 RabbitMQ Server。两个安装包的版本要匹配否则启动会报错或者某些功能不可用。Windows 上安装完成之后进入 RabbitMQ 安装目录的sbin文件夹打开命令行执行rabbitmq-plugins enable rabbitmq_management这个管理界面插件非常重要它能让你在浏览器里看到队列数量、消费速度、消息积压情况。启动完成后访问http://localhost:15672默认账号guest密码guest注意 guest 只能在 localhost 登录如果你需要远程访问管理界面要新建账号并赋予权限。Spring Boot 连接 RabbitMQ 需要知道三个信息地址、端口、虚拟主机。默认的虚拟主机是/用户名密码都是guest。在本 Demo 里我建议新建一个专用虚拟主机比如/es_sync这样避免和系统其他队列混在一起。创建方式可以在管理界面里操作也可以用命令rabbitmqctl add_vhost es_sync rabbitmqctl set_permissions -p es_sync guest .* .* .*2.3 MySQL 初始化脚本MySQL 这边不需要太多花哨操作建一张业务表即可我设计一个简化版的商品表CREATE DATABASE IF NOT EXISTS es_demo DEFAULT CHARACTER SET utf8mb4; USE es_demo; CREATE TABLE product ( id BIGINT(20) NOT NULL AUTO_INCREMENT, name VARCHAR(128) NOT NULL COMMENT 商品名称, category VARCHAR(64) NOT NULL COMMENT 商品分类, price DECIMAL(10,2) NOT NULL COMMENT 价格, description TEXT COMMENT 商品描述, status TINYINT(1) NOT NULL DEFAULT 1 COMMENT 1上架 0下架, create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;插入几条测试数据INSERT INTO product (name, category, price, description) VALUES (Apple MacBook Pro 14, 笔记本电脑, 14999.00, M1 Pro芯片16GB内存512GB固态), (Dell XPS 13, 笔记本电脑, 10999.00, 11代酷睿13.4英寸4K触控屏), (Razer 黑寡妇蜘蛛机械键盘, 外设, 899.00, 绿轴RGB背光全键无冲);3. 核心代码设计与实现3.1 Spring Boot 项目基础配置用 Spring Initializr 创建一个 Spring Boot 2.7.x 项目pom.xml 里重点引入这几个依赖parent groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.14/version /parent dependencies !-- Web 模块 -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency !-- RabbitMQ -- dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency !-- MySQL 驱动和 MyBatis -- dependency groupIdcom.baomidou/groupId artifactIdmybatis-plus-boot-starter/artifactId version3.5.3/version /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-j/artifactId scoperuntime/scope /dependency !-- Elasticsearch Java Client 官方推荐 -- dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version8.3.3/version /dependency dependency groupIdcom.fasterxml.jackson.core/groupId artifactIdjackson-databind/artifactId /dependency /dependencies3.2 Elasticsearch 8.3 的 Java Client 连接方式这里要说一个非常关键的点。很多新手查资料查到 7.x 的RestHighLevelClient然后照抄代码连 8.3结果编译直接报错。ES 官方从 7.15 开始就逐步弃用RestHighLevelClient8.x 全面切换到新的ElasticsearchClient。新客户端的特征是用 Builder 模式创建基于 JSON 请求体构造查询。我提供的连接配置类长这样Configuration public class EsClientConfig { 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(destroyMethod close) public ElasticsearchClient elasticsearchClient() { // ES 8.x 默认开启了 HTTPS本地开发通常关闭或使用证书 // 这里用 http 连接前提是你把 xpack.security.enabled 设为 false // 或者使用 https 且配置 SSL 指纹/证书Demo 里为了简单先用 http final CredentialsProvider credentialsProvider new BasicCredentialsProvider(); credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password)); RestClient restClient RestClient.builder(new HttpHost(host, port, http)) .setHttpClientConfigCallback(httpClientBuilder - httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider)) .build(); ElasticsearchTransport transport new RestClientTransport(restClient, new JacksonJsonpMapper()); return new ElasticsearchClient(transport); } }注意看JacksonJsonpMapper()这一步是用来做请求响应的 JSON 序列化。如果你配置了多个 Jackson ObjectMapper 或者项目里有特定的序列化规则可能会冲突需要单独定制。我在踩坑环节再讲。3.3 实体类和索引 Mapping 设计ES 里的文档数据最好和我们 MySQL 的表结构对应但要考虑 ES 的字段类型设计。同一个字段在 MySQL 里是VARCHAR在 ES 里可能是text分词搜索加keyword精确匹配的组合。我在 Demo 里定义了ProductDocument这个类Setter Getter public class ProductDocument { private Long id; private String name; private String category; private Double price; private String description; private Integer status; private LocalDateTime updateTime; }对应的索引product的创建逻辑public void createProductIndex() throws IOException { ElasticsearchClient client elasticsearchClient(); // 先检查索引是否存在 boolean exists client.indices().exists(e - e.index(product)).value(); if (exists) { return; } client.indices().create(c - c .index(product) .mappings(m - m .properties(id, p - p.long_(l - l)) .properties(name, p - p.text(t - t .fields(keyword, f - f.keyword(k - k)) )) .properties(category, p - p.keyword(k - k)) .properties(price, p - p.double_(d - d)) .properties(description, p - p.text(t - t)) .properties(status, p - p.integer(i - i)) .properties(updateTime, p - p.date(d - d)) ) ); }name字段我设置了text加keyword子字段这个设计在搜索场景里非常实用全文检索时用name精确筛选或排序时用name.keyword。3.4 RabbitMQ 生产者把数据库变更丢进队列业务操作 MySQL 后我们要发一个消息到 RabbitMQ。消息结构我设计成「事件类型 数据 ID」的组合而不是直接发完整数据。理由很简单消费端可能需要回查数据库拿最新数据而且如果消息体太大MQ 的吞吐会受影响。定义消息实体Data public class ProductChangeMessage { /** * 事件类型CREATE / UPDATE / DELETE */ private String eventType; /** * 商品ID消费端根据这个 ID 去 MySQL 查详情 */ private Long productId; }生产者代码Service public class ProductProducer { Autowired private RabbitTemplate rabbitTemplate; public void sendProductChange(String eventType, Long productId) { ProductChangeMessage message new ProductChangeMessage(); message.setEventType(eventType); message.setProductId(productId); // 指定交换机、路由键、消息体 rabbitTemplate.convertAndSend( RabbitConfig.EXCHANGE_NAME, RabbitConfig.ROUTING_KEY, message ); } }这里我用了convertAndSendSpring 会默认用Jackson2JsonMessageConverter把对象转成 JSON前提是你在 RabbitTemplate 里配置了消息转换器否则默认 JDK 序列化会产生很难阅读的二进制内容而且 Java 对象的序列化跨语言基本不可用。3.5 RabbitMQ 配置类与消费者先看配置类Configuration public class RabbitConfig { public static final String EXCHANGE_NAME es.sync.exchange; public static final String QUEUE_NAME es.sync.product.queue; public static final String ROUTING_KEY es.sync.product; Bean public TopicExchange productExchange() { return new TopicExchange(EXCHANGE_NAME, true, false); } Bean public Queue productQueue() { return new Queue(QUEUE_NAME, true); } Bean public Binding productBinding() { return BindingBuilder.bind(productQueue()) .to(productExchange()) .with(ROUTING_KEY); } Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); } Bean public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate template new RabbitTemplate(connectionFactory); template.setMessageConverter(messageConverter()); return template; } }队列声明为持久化交换机也持久化这样 RabbitMQ 重启后消息不会全部丢失。但注意持久化不等于绝对可靠消息确认机制你的生产环境必须另外设计。消费者是核心逻辑所在Component Slf4j public class ProductSyncConsumer { Autowired private ProductMapper productMapper; Autowired private ProductSearchService productSearchService; RabbitListener(queues RabbitConfig.QUEUE_NAME) public void onProductChange(ProductChangeMessage message) throws IOException { log.info(收到商品变更消息: {}, message); Long productId message.getProductId(); if (productId null) { log.warn(消息中 productId 为空忽略); return; } switch (message.getEventType()) { case CREATE: case UPDATE: syncCreateOrUpdate(productId); break; case DELETE: productSearchService.deleteProduct(productId); break; default: log.warn(未知事件类型: {}, message.getEventType()); } } private void syncCreateOrUpdate(Long productId) throws IOException { // 从 MySQL 查询最新数据 Product product productMapper.selectById(productId); if (product null) { // 数据已经被删除同步删除 ES 里的文档 productSearchService.deleteProduct(productId); return; } // 转换为 ES 文档对象并写入 ProductDocument document convertToDocument(product); productSearchService.saveProduct(document); } private ProductDocument convertToDocument(Product product) { ProductDocument doc new ProductDocument(); doc.setId(product.getId()); doc.setName(product.getName()); doc.setCategory(product.getCategory()); doc.setPrice(product.getPrice().doubleValue()); doc.setDescription(product.getDescription()); doc.setStatus(product.getStatus()); doc.setUpdateTime(product.getUpdateTime()); return doc; } }3.6 搜索服务写入与查询 ES我现在补上ProductSearchService的实现。写入用indexAPI查询用bool组合查询这里把常用操作展示出来Service public class ProductSearchService { Autowired private ElasticsearchClient esClient; private static final String INDEX_NAME product; public void saveProduct(ProductDocument doc) throws IOException { esClient.index(i - i .index(INDEX_NAME) .id(doc.getId().toString()) .document(doc) ); } public void deleteProduct(Long productId) throws IOException { esClient.delete(d - d .index(INDEX_NAME) .id(productId.toString()) ); } public ListProductDocument searchByName(String keyword, int page, int size) throws IOException { SearchResponseProductDocument response esClient.search(s - s .index(INDEX_NAME) .query(q - q .bool(b - b .must(m - m .match(mt - mt .field(name) .query(keyword) ) ) .filter(f - f .term(t - t .field(status) .value(1) ) ) ) ) .from((page - 1) * size) .size(size), ProductDocument.class ); return response.hits().hits().stream() .map(Hit::source) .collect(Collectors.toList()); } }这里我特意加了status1的过滤条件因为下架的商品不应该出现在搜索结果里这是实际项目中很常见的筛选需求。4. 业务接口与触发链路4.1 模拟增删改的 Controller为了演示完整流程我写了一个 Controller模拟对商品进行增删改操作每步操作完成后调用 Producer 发消息RestController RequestMapping(/api/product) public class ProductController { Autowired private ProductMapper productMapper; Autowired private ProductProducer productProducer; PostMapping public Long create(RequestBody Product product) { product.setStatus(1); productMapper.insert(product); // 发送创建消息 productProducer.sendProductChange(CREATE, product.getId()); return product.getId(); } PutMapping(/{id}) public Boolean update(PathVariable Long id, RequestBody Product product) { product.setId(id); int rows productMapper.updateById(product); if (rows 0) { productProducer.sendProductChange(UPDATE, id); return true; } return false; } DeleteMapping(/{id}) public Boolean delete(PathVariable Long id) { int rows productMapper.deleteById(id); if (rows 0) { productProducer.sendProductChange(DELETE, id); return true; } return false; } GetMapping(/search) public ListProductDocument search(RequestParam String keyword, RequestParam(defaultValue 1) int page, RequestParam(defaultValue 10) int size) throws IOException { return productSearchService.searchByName(keyword, page, size); } }4.2 为什么不直接双写 MySQL 和 ES看到这里你可能会想直接在一个事务里先写 MySQL 再写 ES 不是更简单吗为什么非要加个 MQ我要说一个很实际的情况。ES 的写入有近实时特性某一瞬间写入 ES 失败比如分片满、网络抖断、集群熔断如果你在主业务事务里同步调用这个失败会直接影响用户操作——下单、改签、发布文章都失败了这在业务里是不可接受的。另一个痛点是如果以后需要把数据同步到 Redis、ClickHouse、或者消息中心每接一个下游就要改一次业务代码改到你怀疑人生。MQ 方案把「MySQL 完成操作」和「ES 更新」解耦业务代码只需要保证消息发出去了消费者负责把消息最终变成搜索引擎里的数据。RabbitMQ 的交换机队列都有确认机制消息不会丢失失败时配合重试和定时补偿能达到最终一致性。这套模式在互联网公司里也是主流做法。5. 常见问题与避坑指南5.1 Elasticsearch 8.3 连接失败的高频原因第一证书校验错误。默认 8.x 开启了 HTTPS 和证书认证如果你用 HTTP 连接会报 SSL 错误。我在本 Demo 里把xpack.security.enabled改成false来简化流程仅限本地开发。生产环境强烈建议保留认证Java Client 侧通过配置 SSL 指纹或者加载证书库来连接代码略显复杂但安全级别完全不同。第二JDK 版本不对。我自己遇到过一次项目用 JDK 8 连接 ES 8.3报java.lang.RuntimeException: java.lang.IllegalStateException: Failed to load a JOSE JWS implementation这个和客户端库的依赖有关ES Java Client 8.x 的elasticsearch-java依赖 Jackson 2.17 以上版本如果项目被 Spring Boot 2.7 锁定到 2.13.5就会冲突。解决方法是至少把jackson-databind升级到 2.14 以上。5.2 RabbitMQ 消息消费不到或消费报错排查队列绑定了但没有消费者监听检查RabbitListener所在的类有没有被 Spring 扫描到队列名是否和配置类一致。消息反序列化失败如果你没配置Jackson2JsonMessageConverter消息内容默认走 JDK 序列化消费者端用ProductChangeMessage接收会报ClassCastException。消息被错误确认丢失消费者处理抛出异常时消息默认会被确认掉。建议在application.yml里设置spring: rabbitmq: listener: simple: acknowledge-mode: auto retry: enabled: true max-attempts: 3 initial-interval: 1000这样消费失败后 RabbitMQ 会不断重试配合死信队列做兜底不至于静默丢消息。5.3 全量数据同步方案与增量同步的配合这个 Demo 演示的是增量同步——MySQL 发生变更才通知 ES 更新。但你第一次上线时MySQL 里已经有一批历史数据ES 索引还是空的。这时需要写一个全量同步接口public void fullSync() throws IOException { ListProduct productList productMapper.selectList(null); for (Product product : productList) { ProductDocument doc convertToDocument(product); productSearchService.saveProduct(doc); } }实际项目中建议先跑全量再开启增量监听确保顺序得当且没有重复数据问题。如果全量量级很大可以分批处理每批 1000 条文档做 bulk 批量写入这样性能数倍提升。5.4 索引字段类型和搜索行为不一致问题ES 的 mapping 一旦创建字段类型就不能在原有的基础上随意修改只能重新建索引再 reindex。所以上线前一定要想清楚字段是text还是keyword。这个性质可以类比 MySQL 的字段类型你不可能同时对一个字段既做模糊匹配又做等值查询除非专门建立前缀索引或生成列。ES 里的text类型就是为全文检索设计的分词后存倒排索引适合match查询keyword类型整存整取适合 term、range、排序、聚合。一个字段又想模糊搜又想精确过滤就设置fields: { keyword: { type: keyword } }子字段两者并存用不同的查询类型区分。我在上面的 mapping 里对name就做了这种处理实际项目里对标题、分类等字段建议都加上。6. 完整 Demo 的扩展方向这套「MySQL RabbitMQ Elasticsearch」的骨架不只是商品搜索能用很多常见业务都能往里面套。比如订单搜索、日志检索、站内文章搜索、会员 CRM 搜索核心链路完全一致只是索引字段和业务事件类型不同。我在实际项目中还加过两个增强一是消息体里带上变更时间戳消费者端做幂等处理防止重复消息导致数据不一致二是增加死信队列消费重试超过 N 次后进入 DLX配合人工补偿或者定时任务扫描补偿。如果你们公司有大数据平台通常还会把 MySQL 的 binlog 通过 Canal 直接投递到 MQ代码层面连侵入都没有这个思路比在业务代码里加 Producer 更解耦但对运维和版本要求更高算是进阶版方案。说到底搜索同步这个需求没有一步到位的银弹重点是要搞明白自己业务里数据一致性要求多高、容忍多大规模的延迟。本地先把这套 Demo 跑通理解消息流转的每一个环节后面上生产再逐步把幂等、重试、监控、全量工具补齐这条路是最稳的。本文还有配套的精品资源点击获取
RELATED — 相关阅读

相关资讯

LATEST — 最新资讯

最新发布

TODAY — 本日精选

新闻

WEEKLY — 本周精选

新闻

MONTHLY — 本月精选

新闻