简介:面向需要打通MySQL与Elasticsearch实时同步链路的Spring Boot开发者,这份Demo完整演示了Elasticsearch 8.3整合、RabbitMQ消息驱动和数据库变更事件处理三大环节,覆盖Maven依赖、连接配置、Repository接口、消息监听、JDBC操作与批量写入等关键实现。压缩包共78个文件,以21个Java源码、21个class文件、17个XML配置和2个YML配置为主,辅以meta、lst、证书及mvnw脚本等,整体约90KB,目录结构清晰便于按模块查阅。通过该项目可理解MySQL→RabbitMQ→Spring Boot→Elasticsearch的数据流转模型,学习异步同步、批量操作、消息重试与错误处理等设计手法,也可直接复用其中的配置与消息处理框架;其中还涉及事件监听和日志记录细节,能帮助读者规避数据不一致与消息丢失等常见问题。目前已有209人学习使用,适合具备一定Spring Boot基础、正在搭建搜索服务或数据同步方案的初中级开发者。
1. 把 MySQL 数据实时同步进 Elasticsearch 8.3:一套 Spring Boot + RabbitMQ 的完整 demo
Spring Boot 整合 Elasticsearch 8.3 并通过 RabbitMQ 同步 MySQL 数据库——这个话题在网上被问的频率很高,但能跑通的完整 demo 很少。原因很简单:ES 8 客户端 API 比 7.x 几乎推翻重写,你照老教程写的RestHighLevelClient一到新版本就全是编译错误;再加上 RabbitMQ 一发一收,序列化对不上消息就直接静默丢失。
这份资源的价值是把「数据库写入 → 消息队列 → 索引更新」串成一个可复现的闭环。它不是只演示某个注解怎么用,而是让你看到一条真实链路:业务写完 MySQL,马上通过 RabbitMQ 通知 ES 更新索引,整个过程在一个 Spring Boot 工程里就能跑起来。
适合正在搭商品搜索、日志检索、订单查询这些功能,又不想用定时任务全量扫表的同学。你先照 demo 跑一遍,再把自己的业务表替换进去,比对着文档从零组装省力得多。
2. 搭好工程骨架:先把 ES 8.3 的依赖版本和连接配置定死
2.1 目录结构与三个组件的职责边界
先说结论:这个 demo 的目录不按“controller、service、dao”这种传统习惯分,而是按组件边界分。我拆完的感受是,把 ES 客户端、RabbitMQ 配置、业务 repository 分开摆放,后续换表、扩队列、同时维护多索引时思路最清晰。
src/main/java/com/example/sync/ ├── SearchSyncApplication.java ├── config/ │ ├── ElasticsearchClientConfig.java │ └── RabbitConfig.java ├── entity/ │ └── Product.java ├── repository/ │ └── ProductRepository.java ├── service/ │ ├── ProductService.java │ └── ProductIndexService.java ├── mq/ │ ├── ProductMessage.java │ ├── ProductSyncProducer.java │ └── ProductSyncConsumer.java └── web/ └── ProductController.javaconfig 包只放两个配置类:一个负责把ElasticsearchClient构建出来扔进 Spring 容器,一个负责声明 RabbitMQ 的交换器、队列、绑定关系。entity 和 repository 是标准的 JPA 写法,Product对应 MySQL 里的业务表。
service 包拆成两个类是有意为之:“业务写库”和“写 ES 索引”是两种不同动作。ProductService只关注 MySQL 层面的增删改,ProductIndexService只关注 ES 的增删改查,互相不掺和。mq 包就是消息链路三个核心:消息体、生产者、消费者。实际项目里如果业务和 ES 索引都不多,把 mq 包直接挂业务模块下没问题;但项目一大,我一般会把消费者单独拆到一个服务里,否则消息一多,业务服务容易跟着重启。从这个目录也能看出来,所有 ES 操作都收在ProductIndexService里,消费者只做「反序列化消息 → 调用索引服务」这一件事,职责很干净。
2.2 pom.xml 依赖怎么配才稳:ES 8 客户端版本不能乱
拆这个 demo 的第一步就是踩版本坑。网上很多文章用spring-boot-starter-data-elasticsearch,它和 ES 8.3 不是同一个维护节奏,版本受 Spring Boot 管控,经常引进来的是 7.x 的客户端。demo 里直接用官方 Java Client,绕开 starter,反而少一层黑匣子。
<parent> <groupId>org.springframework.boot</groupId> <artifactId>spring-boot-starter-parent</artifactId> <version>2.7.18</version> <relativePath/> </parent> <properties> <elasticsearch.version>8.3.3</elasticsearch.version> <java.version>1.8</java.version> </properties>核心依赖就两个,版本必须严格一致:
<dependency> <groupId>co.elastic.clients</groupId> <artifactId>elasticsearch-java</artifactId> <version>${elasticsearch.version}</version> </dependency> <dependency> <groupId>org.elasticsearch.client</groupId> <artifactId>elasticsearch-rest-client</artifactId> <version>${elasticsearch.version}</version> </dependency>elasticsearch-java是 8.x 的官方高级客户端,elasticsearch-rest-client是底层 HTTP 通信的 rest client。它们不像 Spring 的 starter 有版本继承,必须自己写死,而且两个版本号要一模一样。我曾经只升了高版本客户端、忘改 rest-client,启动直接NoClassDefFoundError,查了一个小时才反应过来。
剩下三个依赖用 Spring Boot parent 管理版本:spring-boot-starter-web、spring-boot-starter-data-jpa、spring-boot-starter-amqp,再加一个 MySQL 驱动。JDK 用 1.8 是因为这个 demo 的定位是兼容最广的部署环境;你要是用 JDK 17 也没问题,但建议同时把 Spring Boot 版本提到 2.7 以上的对应版本。
2.3 application.yml:三组连接参数逐项说明
配置集中在application.yml,一共三组连接参数,分开看很清楚:
server: port: 8080 spring: datasource: url: jdbc:mysql://127.0.0.1:3306/search_demo?useUnicode=true&characterEncoding=utf8&serverTimezone=Asia/Shanghai&useSSL=false username: root password: root driver-class-name: com.mysql.cj.jdbc.Driver jpa: hibernate: ddl-auto: update show-sql: true rabbitmq: host: 127.0.0.1 port: 5672 username: guest password: guest virtual-host: / elasticsearch: uris: http://127.0.0.1:9200 username: elastic password: 123456 connect-timeout: 5000 socket-timeout: 30000MySQL 配置里serverTimezone=Asia/Shanghai必须写,不然 MySQL 8 驱动会拿 JVM 默认时区和数据库时区做一次换算,日期字段经常差出 8 小时。JPA 的ddl-auto: update在演示阶段没问题,它会自动建表;正式环境我一般改成validate,避免 Hibernate 误动表结构。
RabbitMQ 的默认账号guest只能在 localhost 下认证,demo 里用127.0.0.1没问题,但你把它换成远程 IP 就会被拒,这是 RabbitMQ 官方策略,不是 bug。Elasticsearch 部分注意password不是默认的123456——ES 8.3 安装后第一次启动会在终端输出一个随机生成的elastic用户密码,你用我写的配置跑的时候,必须把它换成你本地实际安装时生成的那一串,否则第一步就连不上 ES。
3. 构建 Elasticsearch 8.3 客户端:索引设计与基础 CRUD 封装
3.1 用 ElasticsearchClient 直接创建连接:别再把 7.x 的 RestHighLevelClient 拿进来
ES 8 里已经移除了RestHighLevelClient,继续在网上复制 7.x 的写法是纯浪费时间。官方现在给的是elasticsearch-java这套客户端,入口类叫ElasticsearchClient。它底层还是走 rest-client,但所有 API 都改成了函数式链式写法,就拿 connection 工厂来说,代码长这样:
import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.json.jackson.JacksonJsonpMapper; import co.elastic.clients.transport.ElasticsearchTransport; import co.elastic.clients.transport.rest_client.RestClientTransport; import org.apache.http.HttpHost; import org.apache.http.auth.AuthScope; import org.apache.http.auth.UsernamePasswordCredentials; import org.apache.http.client.config.RequestConfig; import org.apache.http.impl.client.BasicCredentialsProvider; import org.elasticsearch.client.RestClient; import org.springframework.beans.factory.annotation.Value; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class ElasticsearchClientConfig { @Value("${elasticsearch.uris}") private String uris; @Value("${elasticsearch.username}") private String username; @Value("${elasticsearch.password}") private String password; @Bean(destroyMethod = "close") public ElasticsearchClient elasticsearchClient() { BasicCredentialsProvider credentialsProvider = new BasicCredentialsProvider(); credentialsProvider.setCredentials(AuthScope.ANY, new UsernamePasswordCredentials(username, password)); RestClient restClient = RestClient.builder(HttpHost.create(uris)) .setHttpClientConfigCallback(builder -> builder .setDefaultCredentialsProvider(credentialsProvider) .setDefaultRequestConfig(RequestConfig.custom() .setConnectTimeout(5000) .setSocketTimeout(30000) .build())) .build(); ElasticsearchTransport transport = new RestClientTransport( restClient, new JacksonJsonpMapper()); return new ElasticsearchClient(transport); } }AuthScope.ANY表示这个账号凭证对所有 host 生效,你不用关心请求打到哪个节点,只要 ES 集群里用户一样就能过认证。connectTimeout是建立 TCP 连接的超时,socketTimeout是等待响应数据的超时——这两个参数在排障时最容易背锅,本地连 ES 报超时不一定是 IP 不通,很可能是 ES 响应慢,socketTimeout设太短被客户端主动掐断。
destroyMethod = "close"保证 Spring 容器关闭时释放底层连接资源。JacksonJsonpMapper是客户端自带的 JSON 映射器,它内部用 Jackson 序列化文档对象。这里有个隐藏细节:默认JacksonJsonpMapper遇到LocalDateTime时的序列化结果是数组[2023,8,1,12,30,0]这种形式,ES date 字段不认识,后面出问题你会很难受,我在避坑章里会专门说解法。
3.2 启动时幂等创建索引:mapping 与 settings 怎么定
索引必须提前存在,否则消费者第一次写数据时 ES 会抛index_not_found_exception。demo 用一个@PostConstruct在 Spring 启动后兜底检查一次索引是否存在,不存在就创建。
import co.elastic.clients.elasticsearch.ElasticsearchClient; import jakarta.annotation.PostConstruct; import org.springframework.stereotype.Service; import java.io.IOException; @Service public class ProductIndexService { private static final String INDEX = "product_index"; private final ElasticsearchClient esClient; public ProductIndexService(ElasticsearchClient esClient) { this.esClient = esClient; } @PostConstruct public void initIndex() throws IOException { boolean exists = esClient.indices() .exists(e -> e.index(INDEX)) .value(); if (!exists) { esClient.indices().create(c -> c .index(INDEX) .mappings(m -> m .properties("id", p -> p.integer()) .properties("name", p -> p.text(t -> t.analyzer("standard"))) .properties("description", p -> p.text(t -> t.analyzer("standard"))) .properties("price", p -> p.double_()) .properties("category", p -> p.keyword()) .properties("createTime", p -> p.date(d -> d.format( "yyyy-MM-dd HH:mm:ss||yyyy-MM-dd||epoch_millis")))) .settings(s -> s .numberOfShards(1) .numberOfReplicas(0))); } } }mapping 里最容易犯的错,是给所有字段都用text。category这类字段适合keyword,因为它是精确匹配、需要聚合统计;name和description才用text做全文检索。analyzer("standard")是 ES 自带分词器,如果服务器上装了 IK 分词插件,你可以直接替换成"ik_max_word",中文搜索效果会好很多,但要记得先确认插件版本和 ES 8.3 匹配,否则索引创建会失败。
numberOfShards(1)和numberOfReplicas(0)是单机演示的低配设定。生产环境至少 3 分片、1 副本,否则节点一挂数据就全没了。这个initIndex的幂等逻辑很重要:如果每次都强行create,第二次启动时应用直接启动失败,exists判断就是给这个兜底。
3.3 封装 save、delete、search:消费者端只认这三个方法
索引服务里所有 ES 操作都集中在一个类里,消费端调用时不需要知道 ES 的 DSL 细节。这是 demo 里最值得抄的部分。
import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch.core.SearchResponse; import co.elastic.clients.elasticsearch.core.search.Hit; import com.example.sync.entity.Product; import org.springframework.stereotype.Service; import java.io.IOException; import java.util.Arrays; import java.util.List; import java.util.stream.Collectors; @Service public class ProductIndexService { private static final String INDEX = "product_index"; private final ElasticsearchClient esClient; public ProductIndexService(ElasticsearchClient esClient) { this.esClient = esClient; } public void save(Product product) throws IOException { esClient.index(i -> i .index(INDEX) .id(String.valueOf(product.getId())) .document(product)); } public void delete(String id) throws IOException { esClient.delete(d -> d .index(INDEX) .id(id)); } public List<Product> searchByKeyword(String keyword, int page, int size) throws IOException { SearchResponse<Product> response = esClient.search(s -> s .index(INDEX) .query(q -> q.multiMatch(m -> m .query(keyword) .fields(Arrays.asList("name", "description")))) .from((page - 1) * size) .size(size), Product.class); return response.hits().hits().stream() .map(Hit::source) .collect(Collectors.toList()); } }index操作对应 MySQL 的 insert 和 update 二合一:ES 是按id覆盖写入的,同一个 id 写两次不会产生重复文档,后写的直接覆盖前一个。这就是同步链路天然幂等的基础,消息重复投递也不会造成脏数据。delete按 id 删文档,删除时 id 不存在也不会报错,ES 会当作“已删除”处理。
searchByKeyword里的multiMatch同时搜了name和description两个字段,比match单字段实用。分页参数from和size注意 ES 默认from+size不能超过 10000,超了报Result window is too large,这个 demo 里没处理,你接真实业务时要么限定深度分页,要么改成search_after。
Product实体要能被 ES 客户端正常序列化,必须满足三个条件:有默认构造函数、有 getter/setter、字段名和 ES mapping 里的 properties 对应。这个实体被 JPA 和 ES 客户端两边共用,字段保持一致,没有额外注解需求,省事很多。
4. 打通 RabbitMQ 链路:先写库再发消息,消费者端同步进 ES
4.1 声明交换器、队列和绑定关系:顺序不重要,类型要对齐
RabbitMQ 的核心是「生产者发到交换器,交换器按路由键把消息投进队列」。demo 里用 DirectExchange,因为它简单:队列和交换器之间通过完全匹配的路由键绑定。配置类长这样:
import org.springframework.amqp.core.Binding; import org.springframework.amqp.core.BindingBuilder; import org.springframework.amqp.core.DirectExchange; import org.springframework.amqp.core.Queue; import org.springframework.amqp.core.QueueBuilder; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class RabbitConfig { public static final String EXCHANGE = "product.sync.exchange"; public static final String QUEUE = "product.sync.queue"; public static final String ROUTING_KEY = "product.sync"; @Bean public DirectExchange productExchange() { return new DirectExchange(EXCHANGE); } @Bean public Queue productQueue() { return QueueBuilder.durable(QUEUE).build(); } @Bean public Binding productBinding() { return BindingBuilder.bind(productQueue()) .to(productExchange()) .with(ROUTING_KEY); } @Bean public MessageConverter jsonMessageConverter() { return new Jackson2JsonMessageConverter(); } }QueueBuilder.durable表示队列持久化,RabbitMQ 重启后队列不会丢。消息本身的持久化还要靠发送端把消息标记为 persistent,Spring AMQP 的convertAndSend默认就是这么做的,所以这里够用。绑定关系声明成@Bean,Spring Boot 启动时会自动执行绑定,不需要额外代码去手动queueBind。
MessageConverter这个 Bean 一定要加。不加的时候,RabbitTemplate默认用SimpleMessageConverter,走 Java 原生序列化,你发送ProductMessage对象会在消息头里带一串 Java 类信息,消费端反序列化容易出兼容性问题;加了这个 Bean 后,生产端和消费端都会自动改成 JSON 格式,消息在 RabbitMQ 管理后台里也能直接看到内容,排查问题方便很多。
4.2 业务服务先写库再发消息:顺序不要反过来
业务层做的是“先保证 MySQL 数据落库,再发同步消息”。这个顺序是同步链路的命门,反过来的话,ES 会先收到消息,而 MySQL 事务还没提交,万一事务回滚了,ES 里就出现一条数据库里根本不存在的记录,索引变成脏数据。
import com.example.sync.entity.Product; import com.example.sync.mq.ProductSyncProducer; import com.example.sync.repository.ProductRepository; import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; @Service public class ProductService { private final ProductRepository productRepository; private final ProductSyncProducer producer; public ProductService(ProductRepository productRepository, ProductSyncProducer producer) { this.productRepository = productRepository; this.producer = producer; } @Transactional public Product save(Product product) { Product saved = productRepository.save(product); producer.sendSave(saved); return saved; } @Transactional public void delete(Integer id) { productRepository.deleteById(id); producer.sendDelete(id); } }其实严格意义上,@Transactional里的save提交后方法才返回,sendSave在事务内执行,但消息发出时数据库事务可能还没真正 commit。常见做法是用TransactionSynchronizationManager.registerSynchronization,在afterCommit回调里发消息,确保 MySQL 提交成功后才通知 ES。demo 为了读代码不绕,把发消息放在事务方法里,这个边界你心里有数就行。生产环境我更加偏向 afterCommit 或 outbox 模式,后者对强一致性要求高,但代码量会增加一个数量级。
消息体ProductMessage设计得尽量简单:
public class ProductMessage { private String operation; // save 或 delete private Integer id; private Product data; public ProductMessage() { } // getters and setters }operation字段区分新增/更新还是删除。保存消息需要带完整data,删除消息只需要id,data为 null 也行。这样消费者不用去 MySQL 回查数据,减少一次数据库访问;也避免了删除时查不到历史数据、压根不知道要删 ES 里哪条记录的尴尬。
生产者类也很直白:
import com.example.sync.entity.Product; import com.example.sync.config.RabbitConfig; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Service; @Service public class ProductSyncProducer { private final RabbitTemplate rabbitTemplate; public ProductSyncProducer(RabbitTemplate rabbitTemplate) { this.rabbitTemplate = rabbitTemplate; } public void sendSave(Product product) { ProductMessage message = new ProductMessage(); message.setOperation("save"); message.setId(product.getId()); message.setData(product); rabbitTemplate.convertAndSend( RabbitConfig.EXCHANGE, RabbitConfig.ROUTING_KEY, message); } public void sendDelete(Integer id) { ProductMessage message = new ProductMessage(); message.setOperation("delete"); message.setId(id); rabbitTemplate.convertAndSend( RabbitConfig.EXCHANGE, RabbitConfig.ROUTING_KEY, message); } }convertAndSend第 3 个参数传对象,Spring 会结合前面配置的Jackson2JsonMessageConverter把对象序列化成 JSON 字符串投递。这里要注意:如果队列和交换器还没绑定好,消息会直接流进 RabbitMQ 默认交换器并“迷路”,所以绑定的@Bean声明是必要条件。
4.3 消费端反序列化后调用索引服务:ack 与异常处理
消费者用@RabbitListener监听队列,拿到消息后根据operation决定是写索引还是删索引。
import com.example.sync.config.RabbitConfig; import com.example.sync.service.ProductIndexService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.amqp.AmqpRejectAndDontRequeueException; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.stereotype.Component; @Component public class ProductSyncConsumer { private static final Logger log = LoggerFactory.getLogger(ProductSyncConsumer.class); private final ProductIndexService productIndexService; public ProductSyncConsumer(ProductIndexService productIndexService) { this.productIndexService = productIndexService; } @RabbitListener(queues = RabbitConfig.QUEUE) public void onMessage(ProductMessage message) { log.info("收到同步消息:operation={},id={}", message.getOperation(), message.getId()); try { if ("save".equals(message.getOperation())) { productIndexService.save(message.getData()); } else if ("delete".equals(message.getOperation())) { productIndexService.delete(message.getId().toString()); } } catch (Exception e) { log.error("写 ES 失败,id={}", message.getId(), e); throw new AmqpRejectAndDontRequeueException(e); } } }这里有两个关键点。第一,@RabbitListener默认采用 AUTO 确认模式:方法正常返回就自动 ack,消息从队列移除;方法抛异常就不 ack,消息重新排队。第二,AmqpRejectAndDontRequeueException的作用是明确告诉监听容器“这条消息别重新入队了”——如果不抛这个异常而直接抛RuntimeException,消息会被放进队列头部反复消费,前台看着队列一直有消息,后台日志疯狂刷错误,这就是典型的消息死循环。
如果你不想让失败消息直接丢,正确做法是给队列配置死信交换器,消费失败 N 次后把消息转到专门的product.sync.dlq,之后再人工或者定时任务去补。demo 没做这一步,但你应该知道这个扩展点在哪。
5. 避坑排查:同步链路里最容易翻车的五个位置
这部分是拆这个 demo 时最扎心的五个位置,每一条我都按「现象 → 原因 → 解决」梳理过,本地跑的时候直接对照定位。
5.1 现象:ES 连接一直报 401 或握手超时
启动应用后日志出现security_exception: missing authentication credentials,或者干脆是Connection timed out。前者基本是 ES 8.3 默认开启 xpack.security,你没带账号或密码不正确;后者常常是连接地址写成了https,但本机 ES 没开 SSL。
解决方式分两步。先确认elasticsearch.yml里的xpack.security.enabled状态和安装时终端打印的 elastic 用户随机密码。再用 http 地址、配置正确的用户名密码;如果 ES 开了xpack.security.http.ssl.enabled,那你还要在RestClient上配setSSLContext,或者更简单地在开发环境先关闭 http.ssl。ES 8 安装后默认密码泄露导致连不上,是做这个 demo 遇到的第一道关卡。
5.2 现象:启动时报 NoClassDefFoundError,是客户端版本不齐
最典型的是NoClassDefFoundError: co/elastic/clients/transport/ElasticsearchTransport,应用连启动都过不去。
原因几乎都是两个依赖版本没有对齐:elasticsearch-java用了 8.3.3,elasticsearch-rest-client还停留在 8.2.x,两个 jar 里的类互相调用时缺口就暴露了。另一种可能是 classpath 里有旧版elasticsearch-rest-high-level-client,它是 7.x 时代的残留物,会和新客户端产生大量冲突。
解决方式:先mvn dependency:tree看实际依赖树,确认没有多余的 es 客户端;然后强制把elasticsearch-java和elasticsearch-rest-client版本统一,放到<properties>里用同一个${elasticsearch.version}引用,从源头杜绝手写版本不一致。
5.3 现象:队列有消息,ES 索引却一条都没增加
RabbitMQ 管理后台能看到队列里有消息,消费者也打印了“收到消息”,但去 ES 查询_count永远是 0。
最常踩的坑是消费端异常被吞了——有的写法是在 catch 里打印日志后直接 return,消息被自动 ack 掉,ES 没写入但你已经无法重试。其次是 ack 策略问题:如果spring.rabbitmq.listener.simple.acknowledge-mode被设成manual,消费端代码里又没手动basicAck,消息会一直处于 Unacked 状态,业务上等于卡死。
解决方式:消费端异常必须向外抛,不要自己吞掉;ack 模式保持auto,让 Spring 根据方法是否抛异常来决定确认还是拒绝。排查的时候先看 RabbitMQ 管理后台里Unacked数字是不是在涨,再看消费者方法 catch 块里有没有log.error,这两处是定位问题的关键。
5.4 现象:LocalDateTime 字段反序列化直接失败
消费者接收消息时报InvalidFormatException: Cannot deserialize value of type java.time.LocalDateTime from String,消息消费不动,队列堆积。
根本原因是 JSON 消息转换器用的ObjectMapper不认识JavaTimeModule。LocalDateTime在没有注册模块时会被序列化成[2023,8,1,10,30,0]数组,而不是字符串,消费端解析当然失败。
解决方式是自定义一个带JavaTimeModule的 ObjectMapper,注册成MessageConverter的底层映射器:
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.SerializationFeature; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @Configuration public class JacksonConfig { @Bean public MessageConverter jsonMessageConverter() { ObjectMapper objectMapper = new ObjectMapper(); objectMapper.registerModule(new JavaTimeModule()); objectMapper.disable(SerializationFeature.WRITE_DATES_AS_TIMESTAMPS); return new Jackson2JsonMessageConverter(objectMapper); } }disable(WRITE_DATES_AS_TIMESTAMPS)后,LocalDateTime会按 ISO 格式序列化,比如2023-08-01T10:30:00。这种格式对 RabbitMQ 本身没问题,但对 ES 的 date 字段来说可能不够友好,所以下一步的日期格式统一也很重要。
5.5 现象:日期写入 ES 后偏移 8 小时或解析报错
ES 索引里有数据了,但时间看起来不对,比如说createTime显示为2023-08-01T20:30:00,而业务库存的是 12:30;或者写入时直接报Can't parse date [2023-08-01T10:30:00]。
原因有两个。一是 ES date 类型内部按 UTC 存储,你传的字符串没有时区标记时,ES 会按 UTC 解析,再显示时又按 UTC 输出,东八区用户看到的就往后偏 8 小时。二是 ES 默认 date 格式只认strict_date_optional_time或epoch_millis其中一种,不带时区的 ISO 格式可能正好不在它接受范围内。
解决方式:在 mapping 里给 date 字段声明多格式,像 demo 里那样写成"yyyy-MM-dd HH:mm:ss||yyyy-MM-dd||epoch_millis",让 ES 至少能识别你传进去的字符串。业务侧则统一在序列化时用@JsonFormat(pattern = "yyyy-MM-dd HH:mm:ss", timezone = "GMT+8")标记LocalDateTime字段,保证 MySQL 读出来的时间和 ES 写入时间一致。这个坑在同步场景里最容易出现,因为 MySQL、RabbitMQ、ES 三层各有一层时间处理,任何一层默认时区不同,最终索引里就差 8 小时。
6. 验收与进阶:Kibana 查询同步结果,再用 bulk 批量补索引
6.1 用 Kibana Dev Tools 走一遍验证命令
同步链路搭好之后,验证不是看“接口没报错”就行,必须沿着数据路径逐层确认。我通常按这个顺序走:
打开 Kibana Dev Tools,先确认索引存在、文档数量正确:
GET /product_index/_count再跑一次关键词搜索,验证 mapping 和 analyzer 配置是否符合预期:
GET /product_index/_search { "query": { "multi_match": { "query": "手机", "fields": ["name", "description"] } } }返回的hits.total应该等于你往 MySQL 里插入、并且消息已消费成功的记录数。如果这个数字对不上,回 RabbitMQ 管理后台看product.sync.queue的Unacked和Ready两个数值,Ready 大于 0 说明消费端没有被唤醒,Unacked 持续增长说明消费者在处理中反复失败,这两个信号能快速定位是链路上哪一段出了问题。
验证时还要注意场景完整性:跑一次新增、跑一次更新、跑一次删除。删除场景尤其容易被忽略,很多人测完新增就认为同步没问题,直到上线后删了 MySQL 记录、ES 里旧数据还在,才意识到 delete 消息路径根本没测。
6.2 批量补数据的 bulk 写法与消费端幂等
如果业务需要一次性导入大量历史数据,一条条index的性能很差。ES 官方建议用 bulk 批量接口,一次请求塞几百条文档。demo 里的索引服务再扩展一个批量方法就行:
import co.elastic.clients.elasticsearch.ElasticsearchClient; import co.elastic.clients.elasticsearch.core.BulkRequest; import co.elastic.clients.elasticsearch.core.BulkResponse; import com.example.sync.entity.Product; import java.io.IOException; import java.util.List; public void bulkSave(List<Product> products) throws IOException { BulkRequest.Builder bulkRequest = new BulkRequest.Builder(); for (Product product : products) { bulkRequest.operations(op -> op.index(idx -> idx .index(INDEX) .id(String.valueOf(product.getId())) .document(product))); } BulkResponse response = esClient.bulk(bulkRequest.build()); if (response.errors()) { // 根据 response.items() 里失败的记录做重试 } }BulkRequest.Builder是典型 builder 模式,循环里每调一次operations就添加一个子请求。这里真正要理解的是:bulk 是幂等的,因为每个子操作都带id,ES 按 id 覆盖写入。这意味着你批量导入过程中如果出现部分失败,重跑一遍不会产生重复文档,已成功的记录被再次覆盖也毫无副作用。这也是整个 MySQL → RabbitMQ → ES 同步链路能稳定运行的基础——重复消费不可怕,可怕的是没有幂等机制导致索引里有大量重复数据。
我现在养成了一个习惯:每配置好一套 MySQL 到 ES 的同步链路,上线前一定强制走一遍三步验证——往业务表写一条测试数据、观察 RabbitMQ 控制台消息是否消费、去 Kibana 查_count和内容,三步全部符合预期才放心交付。同步这种链路,最怕的不是出错,而是出错时你看不到数据在哪一层停下。希望这套 demo 能帮你把每一步都落到看得见的地方,少走点弯路。
本文还有配套的精品资源,点击获取