Spring WebFlux实现Elasticsearch滚动查询同步写入Redis优化
解决Spring WebFlux中Elasticsearch Scroll读取快于Redis写入导致的OOM问题
你的核心问题在于没有让Reactive流的执行顺序对齐:Elasticsearch Scroll请求没有等待Redis写入完成就继续发起,导致内存中积压大量未处理数据。原代码的几个关键问题:
- 使用阻塞式
while循环,违背Reactive非阻塞设计,线程被卡死的同时异步请求不断堆积 - 多处直接调用
subscribe(),把异步操作从主流程剥离,无法保证后续操作等待任务完成 - 批次处理与下一次Scroll请求无依赖关系,ES持续吐数据、Redis写入跟不上,最终引发OOM
修正后的实现方案
通过Reactive链式调用+递归逻辑,严格保证「处理完当前批次Redis写入 → 发起下一次Scroll请求」的顺序,彻底避免数据积压:
// 初始搜索请求(假设已构造好searchRequest,此处省略) CompletableFuture<SearchResponse<EVENT>> initialSearchResponse = elasticsearchAsyncClient.search(searchRequest, EVENT.class); return Mono.fromFuture(initialSearchResponse) .flatMapMany(initialRes -> { long total = initialRes.hits().total().value(); // 先处理第一批数据,再启动递归Scroll流程 return processBatch(initialRes.hits().hits().stream().map(Hit::source).collect(Collectors.toList())) .thenMany(scrollRecursively(initialRes.scrollId(), total, initialRes.hits().hits().size())); }) .then(Mono.just("***********crawlId_cache_completed********* with totalValue processed")); /** * 处理单批次数据:顺序写入Redis,全部完成后才返回 */ private Mono<Void> processBatch(List<EVENT> eventList) { return Flux.fromIterable(eventList) .concatMap(esCache -> { Detail detail = new Detail(esCache.getCrawlId(), LocalDateTime.parse(esCache.getCrawledDate())); List<Detail> details = Collections.singletonList(detail); // 不直接调用subscribe,由concatMap控制执行顺序与等待逻辑 return reactiveRedisOperations.set(esCache.getKey(), details); }) .then(); // 等待当前批次所有Redis写入完成 } /** * 递归执行Scroll:处理完当前批次后,再发起下一次Scroll请求 */ private Flux<Void> scrollRecursively(String scrollId, long total, long processed) { ScrollRequest scrollRequest = new ScrollRequest.Builder() .scrollId(scrollId) .scroll(getScrollTime()) .build(); return Mono.fromFuture(elasticsearchAsyncClient.scroll(scrollRequest, EVENT.class)) .flatMapMany(scrollRes -> { if (scrollRes == null || scrollRes.hits() == null || scrollRes.hits().hits().isEmpty()) { // 无数据时清理Scroll上下文并结束流程 return clearScroll(scrollId).thenMany(Flux.empty()); } long currentBatchSize = scrollRes.hits().hits().size(); long newProcessed = processed + currentBatchSize; log.info("{} of items processed out of {}", newProcessed, total); List<EVENT> eventList = scrollRes.hits().hits().stream().map(Hit::source).collect(Collectors.toList()); // 先完成当前批次Redis写入,再递归执行下一次Scroll return processBatch(eventList) .thenMany(scrollRecursively(scrollRes.scrollId(), total, newProcessed)); }) .onErrorResume(e -> { log.error("Scroll failed: {}", e.getMessage(), e); // 异常时仍需清理Scroll上下文 return clearScroll(scrollId).thenMany(Flux.error(e)); }); } /** * 清理Elasticsearch Scroll上下文的工具方法 */ private Mono<Void> clearScroll(String scrollId) { ClearScrollRequest clearRequest = new ClearScrollRequest.Builder().scrollId(scrollId).build(); return Mono.fromFuture(elasticsearchAsyncClient.clearScroll(clearRequest)) .doOnNext(clearRes -> log.info("Scroll cleared successfully: {}", clearRes.succeeded())) .then(); }
关键改进点说明
- 递归替代阻塞while:用
scrollRecursively实现非阻塞循环,只有当前批次Redis写入完成,才会发起下一次Scroll请求,从根源避免数据积压 - concatMap控制执行顺序:
processBatch中用concatMap保证Redis写入顺序执行,且必须等待当前写入完成才处理下一条数据(若需更高吞吐量,可替换为flatMapSequential) - 统一流生命周期管理:全程不直接调用
subscribe(),所有异步操作纳入Reactive流调度,保证操作间的依赖关系 - 异常与资源清理:添加
onErrorResume保证异常场景下仍能清理Scroll上下文,避免ES端残留无效资源 - 内存友好设计:每批数据处理完成后才加载下一批,内存仅保留当前批次数据,无积压风险
内容的提问来源于stack exchange,提问作者Amulya M
相关产品推荐
相关产品推荐

