ARTICLE DETAIL

资讯详情

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

SpringBoot集成ES:数据同步、异步写入与查询实践

SpringBoot集成ES:数据同步、异步写入与查询实践 先说明一下下面这篇是基于“ES在SpringBoot集成使用”这个项目标题扩展出来的完整技术博文内容围绕标题展开补充了环境准备、客户端选型、配置、查询、数据同步、异步写入、高级特性、问题排查等核心环节并融合了热词涉及的routing、bulk异步写入、minio、mysql同步等常见集成场景。所有内容均基于常见实践补全当作一份个人项目总结来读即可。1. 集成前的环境准备与客户端选型1.1 ES本地安装的两种姿势集成之前先要让ES跑起来。很多人第一次接触ES时会有一个误区直接跑到官网下载最新版本结果SpringBoot项目里用的客户端怎么都连不上其实大概率就是版本不匹配。本地开发我个人推荐用Docker装省事、干净、版本切换方便。一条命令起一个单节点docker run -d \ --name es-single \ -p 9200:9200 -p 9300:9300 \ -e discovery.typesingle-node \ -e xpack.security.enabledfalse \ -e ES_JAVA_OPTS-Xms1g -Xmx1g \ docker.elastic.co/elasticsearch/elasticsearch:7.17.25这里几个参数有讲究discovery.typesingle-node表示单节点模式如果不加ES会尝试做节点发现本地起一个节点时容易进入坑爹的cluster.bootstrapped未完成状态xpack.security.enabledfalse直接关闭安全认证开发环境省去用户名密码的麻烦ES_JAVA_OPTS限制堆内存默认会按物理机大小给笔记本8G内存的话直接给你干到4G其他应用就没法玩了。如果你是Windows本机调试不想用Docker也可以直接下载tar包或zip包。需要注意ES不能使用root用户启动Windows下没有这个问题但Linux/Mac下一定要先创建一个普通用户然后改一下config/elasticsearch.yml里的network.host默认绑定的还是回环地址远程访问会失败。jvm.options里建议把-Xms1g -Xmx1g调一致防止运行中途堆内存动态扩容引发full GC这也是ES生产环境的一个常用优化点。1.2 客户端到底用哪个这是老生常谈但必须明确的选型问题。早期Spring Boot集成ES最常用的方式有三种方案是否推荐理由TransportClient不推荐ES 7.x以后已废弃Java 8之后的项目完全不用考虑RestHighLevelClient7.x可用8.x废弃官方在高版本中标记为deprecated新项目不要再用Spring Data Elasticsearch推荐与Spring Boot集成度高CRUD和简单查询非常顺手Elasticsearch Java Client推荐官方新一代客户端支持异步、版本紧跟ES特性适合写复杂查询我的个人习惯是简单模块用Spring Data Elasticsearch的ElasticsearchRepository复杂的聚合、异步批量、routing操作直接用官方Elasticsearch Java Client两者在同一个项目里是可以并存的并不冲突。版本对应上Spring Boot 2.7.x配ES 7.17是黄金搭档因为Spring Data Elasticsearch 4.4.x对应的就是ES 7.17客户端。Spring Boot 3.x原生对应ES 8.x如果你还在用ES 7.17需要在pom里强制指定elasticsearch-rest-client版本否则会因为版本不一致出现类似NoNodeAvailableException的诡异问题。2. SpringBoot集成ES的核心配置与基础CRUD2.1 依赖引入与yml配置细节以一个Maven项目为例集成ES最核心的依赖就一个dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-elasticsearch/artifactId /dependency如果你的业务需要写一些Native查询或者批量Bulk操作再把官方客户端依赖加上dependency groupIdco.elastic.clients/groupId artifactIdelasticsearch-java/artifactId version7.17.25/version /dependency注意官方客户端在7.17版本时包名是co.elastic.clients在8.x版本也一样但API有变化。我这里以7.17为例是因为大部分存量项目用的还是ES 7.x如果你是新项目直接上8.xAPI里很多类都已经变掉了。yml配置相对简单但有几个参数容易被忽略spring: elasticsearch: uris: - http://localhost:9200 connection-timeout: 3s socket-timeout: 30s username: elastic password: elastic如果是开发环境username和password可以不加。connection-timeout默认只有1秒我个人建议至少给到3秒否则Docker容器刚启动时ES还在恢复分片客户端连一次失败就直接抛异常了。socket-timeout控制的是读取数据的等待时间你如果跑深度分页或者大聚合默认30s可能会超时可以适当调大。2.2 索引映射与实体设计ES和关系型数据库最大的思维差异在于表结构index mapping不是自动从实体推断的它需要你主动设计而且一旦创建成功字段类型基本改不了能做的就是新增字段。我用一个商品索引举例PUT /goods_index { mappings: { properties: { id: { type: keyword }, name: { type: text, analyzer: ik_max_word }, categoryId: { type: keyword }, price: { type: double }, status: { type: integer }, createTime: { type: date, format: yyyy-MM-dd HH:mm:ss||epoch_millis } } } }这里一个最容易翻车的点就是keyword和text的选择。keyword适合做精确查询、排序、聚合比如ID、分类ID、状态码text适合做全文检索存储的原始值会被分词器切词。如果你把商品名称定义成keyword那搜索“手机”就只能精确匹配“手机”搜“苹果手机”完全搜不到反过来如果你把分类ID定义成text做聚合时就会得到一堆奇怪的分词结果。实体类对应如下Data Document(indexName goods_index) public class Goods { Id private String id; Field(type FieldType.Text, analyzer ik_max_word) private String name; Field(type FieldType.Keyword) private String categoryId; Field(type FieldType.Double) private Double price; Field(type FieldType.Integer) private Integer status; Field(type FieldType.Date, format DateFormat.custom, pattern yyyy-MM-dd HH:mm:ss) private LocalDateTime createTime; }Document里的indexName必须是全小写ES索引名不允许大写字母。Id字段映射到ES的_id你要么自己生成一个全局唯一ID比如雪花ID要么让ES自动生成但自动生成的ID是随机的业务查询不方便我一般都自己在应用层赋值。2.3 CRUD操作的几个关键点如果你用了Spring Data ElasticsearchCRUD最省事的姿势是定义Repository接口public interface GoodsRepository extends ElasticsearchRepositoryGoods, String { ListGoods findByStatus(Integer status); PageGoods findByName(String name, Pageable pageable); }单条插入直接用repository.save(goods)批量插入用repository.saveAll(list)。底层实现会自动帮你做bulk不需要自己组织BulkRequest。但注意一个坑如果实体里某个字段是nullSpring Data ES默认也会把该字段一起写入文档覆盖原有文档时会把这个字段直接覆盖成null。如果想要完全忽略null字段实体上可以配合JsonInclude(Include.NON_NULL)这是Jackson的注解Spring Data ES底层用Jackson做序列化所以直接有效。更新文档时需要注意ES的_updateAPI和普通的index请求有本质区别。index是全量覆盖update是部分更新但Spring Data Elasticsearch的save方法底层走的是index所以有时你只想改一个字段结果把其他字段也带过去覆盖了。如果需要部分更新建议直接用官方客户端UpdateRequestString updateRequest new UpdateRequest.Builder() .index(goods_index) .id(id) .doc(Map.of(price, 199.0)) .build();删除操作相对简单repository.deleteById(id)就可以。但要注意ES的删除并不是立刻物理删除它会先标记为已删除后续再触发segment合并时才真正回收空间所以删除后的文档不会立刻在磁盘上消失这是ES正常的工作原理不要误以为代码有问题。3. MySQL与ES的数据同步方案3.1 从业务层面理解“为什么同步难”一个非常现实的问题ES里的数据不会凭空出现多数情况下数据来源是MySQL。并且MySQL是业务系统的权威数据源ES只是查询引擎所以必须保证两边数据最终一致。很多新手第一次做同步时会想在写MySQL的地方顺便再往ES写一份不就行了这种“双写”方案在逻辑上最简单但实际上有严重缺陷如果写ES失败MySQL事务可能已经提交了或者你因为回滚导致MySQL也写失败如果MySQL是主从架构业务代码只连接主库或者只连接从库情况会更复杂。更大的问题是如果业务系统里不止一个服务会修改这张表的数据那每个服务都要维护一段“同时写ES”的逻辑迟早会有漏网之鱼。我在实际项目中见过很多次因为某个角落里的历史数据修复脚本只更新了MySQL没更新ES结果ES和MySQL数据不一致排查起来非常痛苦。所以如果条件允许尽量不要用代码双写而是考虑以下两种更稳妥的方式。3.2 方案一定时任务增量同步这是投入最小、性价比最高的方式适合数据量不大几十万级、对实时性要求不那么极端的场景。核心思路是利用MySQL表中一个update_time字段做增量标识每隔N分钟扫描一次把update_time大于上次同步时间的记录取出来批量写入ES。Component public class GoodsSyncTask { Resource private GoodsMapper goodsMapper; Resource private GoodsRepository goodsRepository; private static final int PAGE_SIZE 500; Scheduled(cron 0 */5 * * * ?) public void syncGoods() { LocalDateTime lastSyncTime getLastSyncTime(); long offset 0; while (true) { ListGoods list goodsMapper.selectByUpdateTimeAfter( lastSyncTime, offset, PAGE_SIZE); if (list.isEmpty()) { break; } goodsRepository.saveAll(list); offset list.size(); if (list.size() PAGE_SIZE) { break; } } } }这个方案有几个必踩的坑。第一个update_time必须有索引否则全表扫描数据量一大直接把MySQL拖垮第二个定时任务每5分钟跑一次建议做分布式锁防止集群多实例同时执行第三个删除同步MySQL里有一行被物理删除了你查不到它ES里就会残留一条脏数据。所以要么用逻辑删除字段deleted代替物理删除要么定期做一次全量对比清理我个人强烈建议用逻辑删除因为逻辑删除天然适合增量同步而且还能查历史数据。3.3 方案二基于Canal订阅binlog如果业务对实时性要求高比如订单状态、物流状态这种需要秒级更新的数据定时任务就不太够用了。常见的做法是引入Canal伪装成MySQL从库订阅binlog变更把变更事件发到消息队列再由消费者写入ES。这套链路是经典的CDCChange Data Capture架构MySQL主库 - Canal - RocketMQ/Kafka - ES消费者。优势是业务代码完全无侵入不需要改任何业务逻辑也不会出现代码双写漏写的问题。缺点是中间组件多运维成本高而且Canal需要目标表使用binlog row格式这是MySQL 8.0的默认配置但如果你的库是老版本还要确认一下。实际落地时消费者写ES一定要用bulk批量处理不能一条一条写。一条update事件写一次ES在高流量下会把ES的写入线程池打满直接出现rejected execution这时候你再看ES日志会看到一堆429异常。3.4 方案三Logstash定时同步低代码消息队列和Canal如果觉得重了Logstash有现成的jdbc插件可以直接定时拉取MySQL数据写入ES。配置一个pipeline文件就行不需要写Java代码。input { jdbc { jdbc_driver_class com.mysql.jdbc.Driver jdbc_connection_string jdbc:mysql://localhost:3306/mall jdbc_user root jdbc_password 123456 statement SELECT * FROM goods WHERE update_time :sql_last_value schedule */5 * * * * * * tracking_column update_time } } output { elasticsearch { hosts [localhost:9200] index goods_index document_id %{id} } }Logstash的方案上手快但要特别小心它的:sql_last_value默认值。如果数据库里最老的数据时间比你设置的上次同步时间还早第一次运行就会漏数据。另外如果你要带中文索引名或字段名注意ID生成规则要显式设定否则每次同步都会生成新文档造成数据翻倍、只增不删的尴尬局面。4. ES异步写入与bulk批量工程化落地4.1 为什么高并发下不能一条条写ES单条并发写入的性能其实并不差但在高并发场景下大量请求到达ES节点后每次CPU都要做分词、索引、刷新segment性能损耗巨大。ES内部每个shard的写入过程大体是请求先进入memory buffer到了refresh间隔再生成一个segment随后segment被commit到磁盘。如果你每秒写几千条文档每条都走完整流程系统开销非常大而且磁盘上会生成大量小segment后续的merge操作会大量消耗IO。bulk批量写入的意义在于一次请求打包几千条文档通过网络层减少交互次数在内存里批量处理刷新segment的次数也大幅降低。实际测试下来相同数据量用bulk写入的吞吐能比单条写入提升5到10倍我调过的项目里商品索引全量重建单条写了700秒改成bulk后90秒就写完了。4.2 Java客户端中的BulkProcessorES 7.17的官方客户端里最常用的批量方式就是BulkProcessor可以设置批量动作阈值、批量大小阈值、刷新间隔、重试策略非常顺手BulkProcessor bulkProcessor BulkProcessor.builder( (request, bulkListener) - request.async().send(bulkListener), new BulkProcessor.Listener() { Override public void beforeBulk(long executionId, BulkRequest request) { // 提交之前可以做日志 } Override public void afterBulk(long executionId, BulkRequest request, BulkResponse response) { // 提交完成后检查hasFailures() } Override public void afterBulk(long executionId, BulkRequest request, Throwable failure) { // 异常处理 } }) .setBulkActions(1000) .setBulkSize(new ByteSizeValue(5, ByteSizeUnit.MB)) .setFlushInterval(TimeValue.timeValueSeconds(5)) .setBackoffPolicy(BackoffPolicy.exponentialBackoff(TimeValue.timeValueMillis(100), 3)) .build();参数的选择从经验来看setBulkActions(1000)意味着攒够1000条就发一次setBulkSize(5MB)是攒够5MB就发一次setFlushInterval(5s)是哪怕没攒够数量5秒也强制发一次避免数据积压在内存中过久。这几个条件谁先触发谁执行。重试策略exponentialBackoff(100ms, 3)的意思是第一次失败后等100毫秒重试第二次等200毫秒第三次等400毫秒最多重试3次。ES在分片正在恢复或者写队列满的时候会返回429这类错误重试是有效的但如果是mapping冲突这种400错误重试也没用直接记日志排查数据才是正解。使用完毕后一定要调用bulkProcessor.close();否则进程退出时内存里可能还有积压的文档没刷出去数据丢失了你还看不出原因。4.3 异步写入的线程池与并发控制代码里request.async().send(bulkListener)本身是异步的但BulkProcessor内部还是会有一个线程去消费你的bulk请求。如果你的应用有多个业务模块同时调同一个BulkProcessor建议把BulkProcessor定义成Spring容器里的单例Bean并且给它指定一个独立的线程池。我自己在项目里踩过的一个坑是BulkProcessor访问ES时的线程池默认用的是EsThreadPoolExecutor提交的任务多了会发生排队积压。一旦应用反复重启旧的请求还会在队列里重试导致新请求无限等待。后来我在调用端统一加了信号量限流比如并发最多50个bulk超过就等一会儿再提交整个写入链路才稳定下来。bulk提交后必须检查response里的hasFailures()。因为一个bulk里混了1000条其中某一条失败了不影响其他条目但如果你不检查就会产生“我以为都写进去了实际上少了几条”的隐患。4.4 routing字段的正确打开方式热搜词里出现了es routing正好展开说一下。routing是ES中控制文档分布到哪个分片的路由参数。默认情况下ES通过hash(_id) % 主分片数来决定文档写入哪个分片如果指定了routing那么就用hash(routing) % 主分片数。这个参数能带来两个实际收益一是提高查询性能如果你查询时带上routingES只会去对应的一个分片查而不是广播到所有分片二是让同一业务维度比如同一个用户、同一个店铺的数据落到同一分片这样后续做聚合或者关联查询时成本更低。写文档时这样设置IndexRequestGoods indexRequest new IndexRequest.BuilderGoods() .index(goods_index) .id(goods.getId()) .routing(goods.getUserId()) .document(goods) .build();读文档时同样带routing可以避免ES在多个分片里挨个找GetRequest getRequest new GetRequest.Builder() .index(goods_index) .id(id) .routing(userId) .build();需要说明的是路由的设计必须在建索引前就想好因为ES的分片数一旦定了就不能动态修改routing规则变更是需要重建索引的大工程。分片数选择上一个索引的分片数不宜过大生产环境里我一般控制在3~9个分片单分片数据量在30GB~50GB以内这种规模下查询和写入的吞吐都是比较均衡的。5. Spring Boot集成ES的高级查询与聚合实操5.1 bool组合查询复杂条件的核心业务中一个常见的场景是“按分类筛选关键字搜索价格区间状态过滤”同时存在。在Spring Data Elasticsearch里你当然可以用方法名推导比如findByCategoryIdAndPriceBetweenAndStatus但这种写法一旦条件超过4个方法名又长又难维护条件动态变化时根本没法写。推荐直接用官方客户端构建bool查询SearchRequest searchRequest new SearchRequest.Builder() .index(goods_index) .query(q - q.bool(b - b .must(m - m.match(mt - mt.field(name).query(手机))) .filter(f - f.term(t - t.field(categoryId).value(c001))) .filter(f - f.range(r - r.field(price).gte(JsonData.of(100)).lte(JsonData.of(5000)))) .filter(f - f.term(t - t.field(status).value(1))) )) .from(0) .size(20) .build();这里最关键的认知是must和filter的区别must里的条件会参与相关性打分影响排序filter里的条件只过滤不计分性能更好。如果只是做约束条件比如“状态1”“分类xx”用filter就对了不要统统塞进must既拖慢查询又把相关性排序搞乱。5.2 高亮查询、分页与深度分页的坑高亮搜索是电商系统的标配。Spring Data ES中可以直接用Native查询构建高亮也可以用官方客户端的HighlightHighlight highlight new Highlight.Builder() .fields(name, new HighlightField.Builder().preTags(span classhl).postTags(/span).build()) .build();这样返回的文档里会多出_highlight字段前端直接渲染。这里有个细节name字段必须是text类型且配置了分词器如果字段是keyword高亮永远不生效。分页方面ES默认的fromsize最多只能到10000条超过之后直接报错。很多人想用from10000,size50翻到第1000页这种思想必须抛弃。ES设计上就不支持传统深分页因为coordinating节点需要把所有分片的数据汇总排序数据量一深内存就爆炸。正确的姿势是search_after游标查询第一次查询时返回的排序字段中取sort数组的值下次查询把它放在search_after参数里不断翻页。这种方案适合“上一页/下一页”的场景但不适合跳到任意一页。如果必须支持任意页码跳转那就只能接受ES 10000条的上限或者把业务拆分到不同程度。5.3 聚合分析从分组到统计一次说清聚合在ES里有三种基础类型桶聚合相当于MySQL的group by、指标聚合相当于count、avg、sum、管道聚合对聚合结果再做聚合。一个实际场景统计每个分类下的商品总数和平均价格。SearchRequest searchRequest new SearchRequest.Builder() .index(goods_index) .size(0) .aggregations(category_count, agg - agg .terms(t - t.field(categoryId).size(10)) .aggregations(avg_price, subAgg - subAgg.avg(a - a.field(price))) ) .build();这里一个最容易犯错的地方字段必须是keyword类型如果字段是textES会报Fielddata is disabled异常。聚合底层用的是doc_valueskeyword和数值类型天然支持text类型直接聚合必须开启fielddata但那是一个非常消耗堆内存的操作生产环境不建议开启。拿到聚合结果后注意terms默认只返回前10个桶size(10)控制桶的数量如果业务需要更多组一定要显式设置size。5.4 一次真实场景的查询代码复盘我最近优化过一个商品搜索接口前端传参与ES条件映射如下前端参数ES查询条件查询类型keywordname字段全文搜索matchcategoryId分类精确匹配termminPrice/maxPrice价格区间rangebrandId品牌IDtermsortField/sortOrder排序sortpage/pageSize分页from/size由于公司业务是后台管理端数据量只有几十万所以直接用了fromsize分页。如果做前台C端搜索数据量过百万我一定会换search_after。每个条件都动态判断是否为空为空就不追加条件。组装好后通过官方客户端发出查询返回结果统一封装成SearchResultVO高亮字段单独取出。这个接口查询耗时从原来走MySQL时的800ms降到了50ms以内。ES在查询层面确实比MySQL在多数复杂搜索场景下更快前提是mapping设计合理、查询类型用得正确。6. 常见问题与排错实录6.1 启动以后“端口连不上”不是ES没起来这个坑出现的频率极高。你检查9200端口确实是通的但Spring Boot应用启动时还是报Connection refused。90%的原因是客户端和ES版本不兼容比如Spring Data Elasticsearch的reactive客户端用的是WebClient底层它连ES 7.x时可能由于请求路径差异遇到兼容问题。排查方式很简单先在命令行里直接执行一次curl http://localhost:9200确认节点信息能返回且version.number是你要的版本再用curl -X GET http://localhost:9200/_cluster/health?pretty看集群状态如果status是yellow说明主分片都可用但副本分片没分配这时写入和查询一般不会有大问题但ES控制台会一直打警告日志对新手来说很容易误判成“ES有问题”。6.2 聚合、高亮、复杂查询报错速查报search_phase_execution_exception通常是跨分片执行时某个分片发生异常根本原因需要看堆栈里的cause。比较常见的是Fielddata is disabled on text fields by default排查方法就是去看mapping确认准备聚合的字段是不是text类型。报IllegalArgumentException: [from] cannot be negative可能是你在代码里给from传了负数或者前端页码传0后你的代码又减1变成了-1。报MapperParsingException通常是写入的字段类型和mapping定义不一致比如mapping里定义status是integer你写入了一段字符串“1”ES会直接拒绝连bulk里的整批数据都会失败需要先把数据清洗干净再写。6.3 磁盘空间不足导致集群只读这个坑很多人项目上线几个月后才会遇到。ES有个默认的cluster.routing.allocation.disk.threshold_enabled配置当节点磁盘使用率超过85%时ES会自动把该节点上的分片迁移走超过90%时会禁止分片分配超过95%时会强制只读此时你写入数据会收到类似disk_watermark_low的报错。一个最直接的临时恢复方法清理掉不再需要的索引或用forcemerge合并segment释放空间等磁盘水位降下来后集群会自动恢复写。千万不要在磁盘满的时候直接重启ES非常容易造成数据损坏。更稳妥的做法是提前规避给每个索引设置生命周期策略比如日志类索引保留7天过期自动删除给ES数据目录单独挂一个大容量磁盘定期用curator或ILM清理旧索引。这些运维习惯在集成阶段就要想好否则后期数据量上来会非常被动。6.4 我的个人排错顺序写ES的过程中debug次数多了我逐渐总结出一套自己的排查顺序第一先看ES集群健康状态如果_cluster/health返回red说明有主分片丢失先恢复集群再说第二看应用日志里具体报的异常类名EsRejectedExecutionException还是ElasticsearchException处理方向完全不同第三看ES端日志在logs/目录里找最近一条error尤其是elasticsearch.log很多时候问题在ES端而不在客户端第四再回头检查Spring Boot配置比如uris是否写错、超时是否太短。在这套顺序下大部分问题十分钟内能定位比遇到问题就翻代码效率高很多。6.5 一个印象深刻的线上事故最后分享一个真实的线上案例。某个订单服务在凌晨大促时突然出现一批订单写入ES失败第一次排查时只看到客户端报BulkRequest has failures但没细看具体失败原因以为是ES集群崩了。后来查ES端日志才发现问题出在订单ID生成重复了导致bulk里出现两条相同_id的文档ES后端抛出version_conflict_engine_exception。这个问题的根源不在ES集成层面而在上游业务系统的ID生成策略。但它给集成ES提了个醒ES的_id必须有唯一性保证千万不能依赖数据库的自增ID因为自增ID在分库分表环境里一定会重复。集成ES前先设计好ID策略这是比任何代码细节都重要的架构决策。ES本身是技术成熟的搜索引擎它和Spring Boot集成时的稳定性瓶颈往往在集成边界上比如数据同步一致性、bulk写入的异常处理、路由策略的合理性。只要把这些边界问题处理妥当ES几乎不会拖后腿剩下的就是把查询条件打磨得更贴合业务罢了。
返回列表