ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

MySQL 实时同步 Elasticsearch 8.3:Spring Boot + RabbitMQ 完整方案

MySQL 实时同步 Elasticsearch 8.3:Spring Boot + RabbitMQ 完整方案 简介面向需要打通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 的完整 demoSpring 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 groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-parent/artifactId version2.7.18/version relativePath/ /parent properties elasticsearch.version8.3.3/elasticsearch.version java.version1.8/java.version /properties核心依赖就两个版本必须严格一致dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version${elasticsearch.version}/version /dependency dependency groupIdorg.elasticsearch.client/groupId artifactIdelasticsearch-rest-client/artifactId version${elasticsearch.version}/version /dependencyelasticsearch-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?useUnicodetruecharacterEncodingutf8serverTimezoneAsia/ShanghaiuseSSLfalse 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 配置里serverTimezoneAsia/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 ListProduct searchByKeyword(String keyword, int page, int size) throws IOException { SearchResponseProduct 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 默认fromsize不能超过 10000超了报Result window is too large这个 demo 里没处理你接真实业务时要么限定深度分页要么改成search_after。Product实体要能被 ES 客户端正常序列化必须满足三个条件有默认构造函数、有 getter/setter、字段名和 ES mapping 里的 properties 对应。这个实体被 JPA 和 ES 客户端两边共用字段保持一致没有额外注解需求省事很多。4. 打通 RabbitMQ 链路先写库再发消息消费者端同步进 ES4.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 重启后队列不会丢。消息本身的持久化还要靠发送端把消息标记为 persistentSpring AMQP 的convertAndSend默认就是这么做的所以这里够用。绑定关系声明成BeanSpring 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删除消息只需要iddata为 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.3elasticsearch-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或者写入时直接报Cant 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 GMT8)标记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(ListProduct 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 是幂等的因为每个子操作都带idES 按 id 覆盖写入。这意味着你批量导入过程中如果出现部分失败重跑一遍不会产生重复文档已成功的记录被再次覆盖也毫无副作用。这也是整个 MySQL → RabbitMQ → ES 同步链路能稳定运行的基础——重复消费不可怕可怕的是没有幂等机制导致索引里有大量重复数据。我现在养成了一个习惯每配置好一套 MySQL 到 ES 的同步链路上线前一定强制走一遍三步验证——往业务表写一条测试数据、观察 RabbitMQ 控制台消息是否消费、去 Kibana 查_count和内容三步全部符合预期才放心交付。同步这种链路最怕的不是出错而是出错时你看不到数据在哪一层停下。希望这套 demo 能帮你把每一步都落到看得见的地方少走点弯路。本文还有配套的精品资源点击获取
返回列表