You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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();
}

关键改进点说明

  1. 递归替代阻塞while:用scrollRecursively实现非阻塞循环,只有当前批次Redis写入完成,才会发起下一次Scroll请求,从根源避免数据积压
  2. concatMap控制执行顺序:processBatch中用concatMap保证Redis写入顺序执行,且必须等待当前写入完成才处理下一条数据(若需更高吞吐量,可替换为flatMapSequential)
  3. 统一流生命周期管理:全程不直接调用subscribe(),所有异步操作纳入Reactive流调度,保证操作间的依赖关系
  4. 异常与资源清理:添加onErrorResume保证异常场景下仍能清理Scroll上下文,避免ES端残留无效资源
  5. 内存友好设计:每批数据处理完成后才加载下一批,内存仅保留当前批次数据,无积压风险

内容的提问来源于stack exchange,提问作者Amulya M

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 08:39:51