☰
SpringBoot整合Elasticsearch 8.3与MySQL同步实战:RabbitMQ实现高效搜索
2026/9/26 5:12:56 网站建设 项目流程

简介:这是一份面向Java后端开发者的实战Demo,聚焦Spring Boot与Elasticsearch 8.3的整合,并借助RabbitMQ打通MySQL到搜索引擎的实时数据同步链路。内容覆盖Spring Data Elasticsearch的Repository抽象、消息监听与发送、数据库变更事件监听、批量写入与错误重试等关键环节,适合正在搭建搜索服务或学习微服务数据流的中高级开发者参考。资源包共78个文件,以21个java源码、21个class编译文件、17个xml配置及2个yml配置文件为主,另含证书、构建脚本与说明文档,压缩包约90KB,结构紧凑便于快速导入IDE运行调试。目前已有209人学习下载。通过该Demo可掌握MySQL→RabbitMQ→Spring Boot→Elasticsearch的完整同步链路,理解异步解耦、批量优化与监控排错思路,对电商、日志分析等实时检索场景有直接借鉴价值。

1. SpringBoot 整合 Elasticsearch 8.3 并同步 MySQL:这套组合到底解决什么问题

电商后台的商品搜索是个典型场景:MySQL 里存着几十万条商品记录,运营改一个价格、下架一个 SKU,用户端搜索必须尽快反映出来。直接拿 MySQL 的LIKE '%关键词%'去扛搜索流量,数据量一上来查询就会明显变慢,字段多了连索引都用不上。把检索交给 Elasticsearch、把 MySQL 当唯一事实源、中间用 RabbitMQ 做变更事件的搬运工,是这几年在 SpringBoot 项目里最常见的一种落地组合。

这篇要讲的就是这条链路怎么从零跑通:SpringBoot 怎么连上 Elasticsearch 8.3、MySQL 的增删改怎么变成消息、消费者怎么把消息写进索引、以及全量初始化与增量同步怎么配合。适合手里有 SpringBoot 项目、想给业务加一层搜索能力、又不想引入 Canal 这类额外组件的后端同学。下面按「先跑通单点、再串链路、最后处理一致性」的顺序展开,每一步都给可复制的配置和代码。

2. 环境搭建:Elasticsearch 8.3、RabbitMQ、MySQL 三件套怎么起

2.1 Elasticsearch 8.3 的启动与安全配置

Elasticsearch 8.x 和 7.x 最大的差别是默认开启了安全认证,第一次启动会生成证书和密码,很多人卡在这一步。用 Docker 起一个单节点最省事,注意内存参数别给太小,否则容器会反复重启。

# 单节点开发环境,关闭不必要的内存占用 docker run -d --name es83 \ -p 9200:9200 -p 9300:9300 \ -e "discovery.type=single-node" \ -e "ES_JAVA_OPTS=-Xms1g -Xmx1g" \ -e "xpack.security.enabled=true" \ docker.elastic.co/elasticsearch/elasticsearch:8.3.3

启动后需要拿到初始密码,执行下面这条命令,它会重置elastic用户的密码并输出:

docker exec -it es83 bin/elasticsearch-reset-password -u elastic

逻辑说明:discovery.type=single-node让节点跳过集群发现,开发机不用配多节点;xpack.security.enabled=true是 8.x 的默认行为,显式写出来是为了提醒你后面连接必须带账号密码。参数上-Xms和-Xmx设成一样能避免堆动态调整带来的停顿,开发机 1g 够用,生产按数据量给到物理内存一半以内。

验证是否起来,用 curl 带 basic auth 访问:

curl -k -u elastic:你重置后的密码 https://localhost:9200

返回带cluster_name和version.number的 JSON 就说明单点通了。注意 8.x 默认走 https,-k是跳过自签证书校验,生产环境要换成正式证书。

2.2 RabbitMQ 的安装与 virtual host 权限坑

RabbitMQ 用 Docker 起同样快,但管理界面能打开、admin 却建不了虚拟主机,是新手最常撞的墙。原因是admin这个账号默认只有管理界面的登录权限,没有对具体 vhost 的配置和读写权限。

docker run -d --name rabbitmq \ -p 5672:5672 -p 15672:15672 \ rabbitmq:3.11-management

起来后进容器创建账号和 vhost:

docker exec -it rabbitmq rabbitmqctl add_user app app123456 docker exec -it rabbitmq rabbitmqctl add_vhost /es-sync # 关键一步:给账号授予 vhost 的完整权限 docker exec -it rabbitmq rabbitmqctl set_permissions -p /es-sync app ".*" ".*" ".*" docker exec -it rabbitmq rabbitmqctl set_user_tags app administrator

逻辑说明:set_permissions的三个正则分别对应 configure、write、read,".*"表示全放开,开发环境够用,生产要按队列名收窄。set_user_tags给 administrator 标签只是让它能进管理界面,和 vhost 权限是两码事,别混。参数上-p /es-sync指定操作哪个虚拟主机,漏了会作用到默认的/。

2.3 MySQL 侧的准备与 binlog 认知

MySQL 这边不需要额外组件,但要想清楚一件事:我们靠应用层发消息,不依赖 binlog 解析。所以建表时给需要同步的表加一个update_time字段,全量初始化时按它排序分页,增量时靠业务代码在写库后发消息。

CREATE TABLE product ( id BIGINT PRIMARY KEY AUTO_INCREMENT, name VARCHAR(128) NOT NULL, price DECIMAL(10,2) DEFAULT 0, status TINYINT DEFAULT 1, update_time DATETIME DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4;

逻辑说明:ON UPDATE CURRENT_TIMESTAMP让每次更新自动刷新时间,全量同步时用它做游标能避免漏数据。status默认给 1 而不是 0,是因为很多业务把 0 当删除标记,默认值设错会导致新数据一进来就被过滤掉,这个坑在热词里也有人提过。

3. SpringBoot 接入 Elasticsearch 8.3:客户端选型与最小可跑代码

3.1 用 Spring Data Elasticsearch 还是官方 Java API Client

SpringBoot 3.x 对应 Spring Data Elasticsearch 5.x,底层已经换成官方的elasticsearch-java客户端。选型上分两种情况:如果只是简单的 CRUD 和少量条件查询,用 Spring Data 的ElasticsearchRepository最省代码;如果要拼复杂的 bool 查询、聚合、自定义打分,直接用官方ElasticsearchClient更灵活。

我一般两个都引,简单索引用 Repository,复杂查询注入ElasticsearchClient。依赖这样加:

<dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-data-elasticsearch</artifactId> </dependency> <dependency> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-amqp</artifactId> </dependency>

逻辑说明:spring-boot-starter-data-elasticsearch会自动带上elasticsearch-java和jackson,不用单独引客户端版本,避免版本冲突。spring-boot-starter-amqp是 RabbitMQ 的封装,后面发消息用RabbitTemplate。

配置文件里把 ES 和 RabbitMQ 的连接写清楚:

spring: elasticsearch: uris: https://localhost:9200 username: elastic password: 你重置后的密码 connection-timeout: 5s rabbitmq: host: localhost port: 5672 username: app password: app123456 virtual-host: /es-sync

逻辑说明:uris用 https,因为 8.x 默认开 TLS;virtual-host必须和前面创建的/es-sync一致,写错会报NOT_ALLOWED。connection-timeout设短一点,ES 挂了能快速失败而不是拖住启动。

3.2 定义索引映射与实体类

索引映射建议手动建,不要全靠自动推断,否则price可能被识别成 text 导致范围查询失效。

curl -k -u elastic:密码 -X PUT "https://localhost:9200/product" -H 'Content-Type: application/json' -d ' { "mappings": { "properties": { "id": { "type": "long" }, "name": { "type": "text", "analyzer": "standard" }, "price": { "type": "scaled_float", "scaling_factor": 100 }, "status": { "type": "integer" }, "updateTime": { "type": "date", "format": "yyyy-MM-dd HH:mm:ss" } } } }'

逻辑说明:name用 text 走分词,方便模糊匹配;price用scaled_float而不是 double,金额场景能省空间又避免浮点误差;updateTime指定格式,和 Java 侧的LocalDateTime对齐。参数scaling_factor: 100表示保留两位小数。

实体类对应写:

@Document(indexName = "product") public class ProductDoc { @Id private Long id; private String name; private BigDecimal price; private Integer status; @Field(type = FieldType.Date, format = DateFormat.date_hour_minute_second) private LocalDateTime updateTime; // getter/setter 省略 }

逻辑说明:@Document指定索引名,@Id对应文档 id,用 MySQL 主键保证幂等。@Field的 format 要和映射里的 format 一致,不一致写入会报解析异常。

3.3 一个能跑通的写入与查询

@Service public class ProductSearchService { @Autowired private ElasticsearchClient client; public void save(ProductDoc doc) throws IOException { client.index(i -> i .index("product") .id(String.valueOf(doc.getId())) .document(doc)); } public List<ProductDoc> search(String keyword) throws IOException { SearchResponse<ProductDoc> resp = client.search(s -> s .index("product") .query(q -> q.match(m -> m.field("name").query(keyword))), ProductDoc.class); return resp.hits().hits().stream() .map(Hit::source).toList(); } }

逻辑说明:index方法用主键做文档 id,重复写入会覆盖,天然幂等,这是后面消息重投不产生脏数据的基础。search用 match 查询走分词,返回直接映射成实体。参数上field("name")要和映射里的字段名一致,写错会返回空结果而不是报错,排查时先看 mapping。

4. RabbitMQ 串起同步链路:从 MySQL 变更到索引更新

4.1 交换机、队列与绑定关系的设计

同步链路的消息模型要提前定好,否则后面加表、加索引会乱。常见做法是用一个 topic 交换机,路由键按「表名.操作」命名,队列按索引维度分。

@Configuration public class RabbitConfig { public static final String EXCHANGE = "db.sync.exchange"; public static final String QUEUE = "product.sync.queue"; @Bean public TopicExchange syncExchange() { return new TopicExchange(EXCHANGE, true, false); } @Bean public Queue productQueue() { return QueueBuilder.durable(QUEUE).build(); } @Bean public Binding productBinding() { return BindingBuilder.bind(productQueue()) .to(syncExchange()).with("product.#"); } }

逻辑说明:TopicExchange的durable=true保证 broker 重启后交换机还在;队列同样持久化。绑定用product.#通配,product.insert、product.update、product.delete都能路由到这个队列。参数上第二个false是 autoDelete,别设 true,否则没有消费者时交换机会被删。

4.2 生产端:写库后发消息的正确姿势

最容易翻车的地方是「先发消息还是先提交事务」。如果先发消息、事务回滚了,索引里就多了一条 MySQL 没有的数据。正确顺序是先提交事务,再发消息。

@Service public class ProductService { @Autowired private ProductMapper mapper; @Autowired private RabbitTemplate rabbitTemplate; @Transactional public void updateProduct(Product p) { mapper.updateById(p); } // 事务提交后再发,用 TransactionSynchronization 保证顺序 public void updateAndSync(Product p) { updateProduct(p); TransactionSynchronizationManager.registerSynchronization( new TransactionSynchronization() { @Override public void afterCommit() { rabbitTemplate.convertAndSend( RabbitConfig.EXCHANGE, "product.update", JSON.toJSONString(p)); } }); } }

逻辑说明:afterCommit回调保证消息只在事务真正提交后发出,避免回滚导致的脏索引。消息体直接发 JSON 字符串,消费端反序列化。参数上路由键product.update要和绑定规则匹配,写错消息会被丢弃且不报错,排查时看 RabbitMQ 管理界面的消息速率。

4.3 消费端:把消息写进 Elasticsearch

@Component public class ProductSyncConsumer { @Autowired private ProductSearchService searchService; @RabbitListener(queues = RabbitConfig.QUEUE) public void onMessage(String body) { Product p = JSON.parseObject(body, Product.class); ProductDoc doc = new ProductDoc(); BeanUtils.copyProperties(p, doc); try { searchService.save(doc); } catch (IOException e) { // 抛出去让 RabbitMQ 走重试或死信 throw new RuntimeException("同步失败 id=" + p.getId(), e); } } }

逻辑说明:消费端用主键覆盖写入,重复消费不会产生重复文档。捕获异常后重新抛出,交给 RabbitMQ 的重试机制处理,不要吞掉异常,否则消息被 ack 但数据没同步,问题会被藏起来。参数上@RabbitListener默认手动 ack 关闭,即自动 ack,生产建议改成手动 ack 配合死信队列。

4.4 全量初始化与增量同步的配合

上线时索引是空的,得先把 MySQL 存量数据灌进去,再开启增量。全量用分页游标,按update_time和id排序,避免深分页。

public void fullSync() { long lastId = 0L; int size = 500; while (true) { List<Product> list = mapper.selectByCursor(lastId, size); if (list.isEmpty()) break; for (Product p : list) { ProductDoc doc = new ProductDoc(); BeanUtils.copyProperties(p, doc); try { searchService.save(doc); } catch (IOException ignored) {} lastId = p.getId(); } } }

逻辑说明:用id > lastId的游标分页代替limit offset,数据量大时不会越翻越慢。每批 500 条是经验值,太大内存吃紧,太小网络往返多。全量跑完后,增量消息可能已经积压在队列里,消费端幂等写入正好把重复的覆盖掉,不用额外去重。

5. 避坑与排查:同步链路上最容易翻车的 5 个点

5.1 现象:ES 健康检查一直失败,应用启动报 health check failed

原因:SpringBoot 的ElasticsearchRestClientHealthIndicator会定期探测 ES,8.x 默认 https 加自签证书,客户端校验证书失败就报健康检查不通过。

解决:开发环境在application.yml里关掉证书校验,或者把自签证书导入信任库。生产环境用正式证书,不要图省事关校验。

spring: elasticsearch: uris: https://localhost:9200 username: elastic password: 密码

如果还报错,先单独用 curl 确认 ES 本身能访问,再排查应用侧配置。

5.2 现象:消息发出去了,索引里却没数据

原因:路由键和绑定规则不匹配,消息被交换机丢弃。topic 交换机在没有匹配队列时默认静默丢弃,不报错。

解决:去 RabbitMQ 管理界面的 Exchanges 页看绑定关系,确认路由键能被product.#匹配。也可以在交换机上加一个备用交换机(alternate-exchange)兜底,把没路由的消息收集起来排查。

5.3 现象:admin 用户登录管理界面正常,但代码连接报 NOT_ALLOWED

原因:账号没有目标 vhost 的权限,或者virtual-host配置写成了/而实际用的是/es-sync。

解决:用rabbitmqctl list_permissions -p /es-sync确认权限,缺了就用set_permissions补。配置里的virtual-host要和创建时完全一致,大小写敏感。

5.4 现象:全量同步跑一半内存溢出

原因:一次性把整表查进内存,或者分页用limit offset在深分页时 MySQL 扫描大量行。

解决:改成游标分页,每批固定条数,处理完一批再查下一批。同时给update_time或id建索引,让游标查询走索引。

5.5 现象:更新操作偶尔丢失,索引数据和 MySQL 对不上

原因:先发消息后提交事务,事务回滚了消息却已经发出;或者消费端吞了异常导致消息被 ack 但没写入。

解决:生产端用afterCommit回调发消息;消费端异常必须抛出,配合手动 ack 和死信队列,让失败的消息能被重投或人工处理。定期跑一次对账任务,比对 MySQL 和 ES 的条数与关键字段,发现偏差就触发补偿。

6. 进阶:用版本号做幂等与对账,把一致性握在手里

消息重投、乱序是分布式同步绕不开的问题。前面靠主键覆盖写入解决了「重复」,但解决不了「乱序」——比如一条 update 消息晚于 delete 消息到达,索引里就会复活一条已删除的数据。我一般会在消息体里带一个版本号,用 MySQL 的update_time毫秒值或者自增版本字段,消费端写入前先比对。

public void saveWithVersion(ProductDoc doc, long version) throws IOException { // 用脚本更新,只有新版本号大于已存版本才写入 client.update(u -> u .index("product") .id(String.valueOf(doc.getId())) .script(s -> s.inline(i -> i .source("if (ctx._source.version == null || params.v > ctx._source.version) { ctx._source = params.doc; ctx._source.version = params.v; }") .params(Map.of("v", version, "doc", doc)))), ProductDoc.class); }

逻辑说明:用 painless 脚本做条件更新,只有版本号更大才覆盖,天然抵御乱序。参数params.v是消息里的版本号,params.doc是新文档内容。注意脚本更新要求文档已存在,首次写入还是走 index。

对账任务可以每天凌晨跑一次,按update_time拉出最近一天变更的 id,逐个查 ES 比对关键字段,不一致的重新发一条同步消息。这个补偿逻辑不复杂,但能兜住绝大多数偶发不一致。

// 对账伪代码:找出 MySQL 有、ES 没有或字段不一致的记录 List<Long> changedIds = mapper.selectChangedSince(lastCheckTime); for (Long id : changedIds) { Product db = mapper.selectById(id); ProductDoc es = searchService.getById(id); if (es == null || !Objects.equals(db.getPrice(), es.getPrice())) { rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE, "product.update", JSON.toJSONString(db)); } }

逻辑说明:对账只处理变更过的 id,不做全表扫描,成本可控。比对字段按业务重要性挑,价格、状态、名称这几个最关键。发现不一致就重发消息,走同一条消费链路,不用另写修复逻辑。

血泪经验是:别指望消息百分百不丢不重,把幂等和对账做扎实,比追求「绝对可靠」的消息投递更实际。我现在的习惯是,任何走消息同步的链路,上线前先把对账脚本写好,出问题时有后悔药可吃。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询