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

如何让Reactive Elasticsearch Scroll查询等待所有数据消费完成?

问题描述

需以响应式方式查询Elasticsearch滚动索引。因@Document注解不支持通过Spring EL动态指定索引名(如@Document(index = "indexName-#{new Date().format(yyyy-MM-dd)}")),故采用ReactiveElasticsearchTemplate实现运行时动态修改索引名。

但数据量超10000条,需用Scroll机制循环查询直至获取全部数据。目前已完成首次查询和滚动查询的代码实现且能返回数据,但需合并所有结果后再返回。当前问题:消费者仍在处理数据时,空结果已返回给前端。

原代码如下:

public Flux<ELKModel> getByTradeDateBetween(LocalDateTime from, LocalDateTime to)
  throws Exception {
List<ELKModel> result = new ArrayList<ELKModel>();
List<Long> total = new ArrayList<>();
List<Long> currentSize = new ArrayList<>();
List<String> scrollId = new ArrayList<>();

NativeSearchQueryBuilder sourceBuilder = new NativeSearchQueryBuilder();
sourceBuilder.withQuery(
    QueryBuilders.boolQuery().must(QueryBuilders.rangeQuery(TRADE_DATE).gte(from).lte(to)));
sourceBuilder.withPageable(PageRequest.of(0, SINGLE_QUERY_SIZE));
NativeSearchQuery query = sourceBuilder.build();
 elasticsearchSupport
    .scrollStart(query, ELKModel.class)
    .map(ELKModelWrapper::valueFrom)
        subscribe(
        wrapper -> {
          total.add(wrapper.getTotal());
          currentSize.add(wrapper.getCurrentSize());
          result.addAll(wrapper.getResults());
          scrollId.add(wrapper.getScrollId());
        }).dispose();

while (currentSize.size() == 1 && total.size() == 1 && currentSize.get(0) < total.get(0)) {
    elasticsearchSupport
      .scrollContinue(scrollId.get(0), ELKModel.class)
      .map(ELKModelWrapper::valueFrom)
      .subscribe(
          wrapper -> {
            currentSize.add(0, currentSize.get(0) + wrapper.getCurrentSize());
            result.addAll(wrapper.getResults());
            scrollId.add(0, wrapper.getScrollId());
          }).dispose();
          
}

return Flux.fromIterable(result);
}
问题原因
  1. 异步与同步逻辑冲突:代码用subscribe()触发异步查询后立即终止订阅,同时用同步while循环判断滚动条件。但异步操作还未完成时,while循环已经执行完毕,此时result集合未被填充,直接返回空的Flux。
  2. 线程安全风险:使用非线程安全的ArrayList在异步回调中操作,可能引发并发修改异常。
  3. Scroll资源未清理:未在查询结束后清理Elasticsearch的Scroll上下文,会造成节点资源泄漏。
解决方案

采用Reactor响应式原生操作符串联Scroll查询流程,避免异步同步混合问题,同时保证线程安全并正确清理资源:

public Flux<ELKModel> getByTradeDateBetween(LocalDateTime from, LocalDateTime to) {
    // 构建查询,必须设置Scroll过期时间(示例为5分钟)
    NativeSearchQuery query = new NativeSearchQueryBuilder()
            .withQuery(QueryBuilders.boolQuery()
                    .must(QueryBuilders.rangeQuery(TRADE_DATE).gte(from).lte(to)))
            .withPageable(PageRequest.of(0, SINGLE_QUERY_SIZE))
            .withScroll(Duration.ofMinutes(5))
            .build();

    // 首次发起Scroll查询,递归滚动直至获取全部数据
    return elasticsearchSupport.scrollStart(query, ELKModel.class)
            .map(ELKModelWrapper::valueFrom)
            // 用expand递归触发下一次滚动查询
            .expand(wrapper -> {
                // 判断是否需要继续滚动:当前已获取数据量小于总条数,且本次返回结果非空
                if (wrapper.getCurrentSize() > 0 && wrapper.getCurrentSize() < wrapper.getTotal()) {
                    return elasticsearchSupport.scrollContinue(wrapper.getScrollId(), ELKModel.class)
                            .map(ELKModelWrapper::valueFrom)
                            // 更新累计获取的数据量
                            .map(nextWrapper -> {
                                nextWrapper.setTotal(wrapper.getTotal());
                                nextWrapper.setCurrentSize(wrapper.getCurrentSize() + nextWrapper.getCurrentSize());
                                return nextWrapper;
                            });
                } else {
                    // 无需继续滚动,返回空Flux终止递归
                    return Flux.empty();
                }
            })
            // 展开所有结果
            .flatMap(wrapper -> Flux.fromIterable(wrapper.getResults()))
            // 合并所有结果后再返回(若无需合并,可删除此行)
            .collectList()
            .flatMapMany(Flux::fromIterable)
            // 最后清理Scroll上下文,避免资源泄漏
            .doFinally(signalType -> {
                // 根据实际实现清理ScrollId,示例:
                // elasticsearchSupport.scrollClear(lastScrollId);
            });
}
关键说明
  • 用expand实现递归滚动:完全遵循响应式模型,避免同步阻塞,自动等待每次异步查询完成后判断是否继续滚动。
  • 线程安全的结果处理:通过Reactor操作符处理结果,无需手动使用线程不安全集合。
  • 强制清理Scroll资源:必须在查询结束后清理ES的Scroll上下文,否则会持续占用节点内存。
  • 设置Scroll过期时间:首次查询时必须指定,ES会自动清理超时的Scroll,主动清理更可靠。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 01:15:02