如何让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); }
问题原因
- 异步与同步逻辑冲突:代码用
subscribe()触发异步查询后立即终止订阅,同时用同步while循环判断滚动条件。但异步操作还未完成时,while循环已经执行完毕,此时result集合未被填充,直接返回空的Flux。 - 线程安全风险:使用非线程安全的
ArrayList在异步回调中操作,可能引发并发修改异常。 - 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
相关产品推荐
相关产品推荐

