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

如何从Flux<AggregationContainer>获取聚合桶计数及Elasticsearch总记录数

Spring WebFlux 操作Elasticsearch聚合:解析桶计数与总记录数

在Spring WebFlux中使用Elasticsearch聚合功能时,需要实现两个核心需求:

  • 从Flux<AggregationContainer<?>>中提取聚合桶的文档计数(即响应中buckets里的doc_count)
  • 获取查询对应的Elasticsearch总记录数

示例请求与响应

请求DSL

POST order/_search
{
  "size": 0, 
  "aggs": {
    "status": {
      "terms": {
        "field": "order_status",
        "size": 100
      }
    }
  }
}

响应示例

{
  "took": 1,
  "timed_out": false,
  "_shards": {
    "total": 1,
    "successful": 1,
    "skipped": 0,
    "failed": 0
  },
  "hits": {
    "total": {
      "value": 37,
      "relation": "eq"
    },
    "max_score": null,
    "hits": []
  },
  "aggregations": {
    "status": {
      "buckets": [
        {
          "key": "CREATED",
          "doc_count": 34
        },
        {
          "key": "CANCELLED",
          "doc_count": 2
        },
        {
          "key": "COMPLETED",
          "doc_count": 1
        }
      ]
    }
  }
}

1. 解析聚合桶的文档计数

ReactiveElasticsearchOperations.aggregate()返回的Flux<AggregationContainer<?>>需先转换为Mono(聚合结果通常为单个容器),再强转为TermsAggregation类型,遍历桶即可提取key和doc_count:

@Autowired
private ReactiveElasticsearchOperations operations;

public Mono<Map<String, Long>> getOrderStatusAggregation() {
    // 构建聚合查询
    SearchQuery searchQuery = new NativeSearchQueryBuilder()
            .withQuery(QueryBuilders.matchAllQuery()) // 替换为实际业务查询条件
            .withSize(0) // 不返回具体文档,仅获取聚合结果
            .addAggregation(AggregationBuilders.terms("status")
                    .field("order_status")
                    .size(100))
            .build();

    return operations.aggregate(searchQuery, "order") // 指定目标索引名
            .single() // 聚合结果仅一个容器,取单个结果
            .map(aggregationContainer -> {
                // 强转为Terms类型聚合
                Terms termsAgg = aggregationContainer.unwrap(Terms.class);
                Map<String, Long> statusCountMap = new HashMap<>();
                // 遍历所有桶,提取状态与对应计数
                for (Terms.Bucket bucket : termsAgg.getBuckets()) {
                    String status = bucket.getKeyAsString();
                    long count = bucket.getDocCount();
                    statusCountMap.put(status, count);
                }
                return statusCountMap;
            });
}

2. 获取总记录数

总记录数可从SearchHits的total字段获取,需使用search()方法执行查询(同时保留聚合配置):

public Mono<Long> getTotalRecordCountWithAggregation() {
    SearchQuery searchQuery = new NativeSearchQueryBuilder()
            .withQuery(QueryBuilders.matchAllQuery())
            .withSize(0)
            .addAggregation(AggregationBuilders.terms("status")
                    .field("order_status")
                    .size(100))
            .build();

    return operations.search(searchQuery, Count.class, "order")
            .map(SearchHits::getTotalHits)
            .map(TotalHits::value);
}

若仅需总记录数、无需聚合结果,可直接使用count()方法:

public Mono<Long> getTotalRecordCount() {
    SearchQuery countQuery = new NativeSearchQueryBuilder()
            .withQuery(QueryBuilders.matchAllQuery())
            .build();
    return operations.count(countQuery, "order");
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 14:20:23