本文目录导读:

- 环境准备
- 案例一:简单的数值统计(最大值、最小值、平均值、总和)
- 案例二:分组统计(Terms聚合)
- 案例三:日期直方图聚合(按时间统计)
- 案例四:嵌套聚合(多维度分析)
- 案例五:百分比统计(Percentiles)
- 实用建议
在Java中使用Elasticsearch的聚合功能进行数据统计,通常通过elasticsearch-rest-high-level-client(7.x版本)或elasticsearch-java(8.x版本)实现。
以下是几个典型的Java ES聚合统计案例:
环境准备
Maven依赖(7.x版本推荐):
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<version>7.17.0</version>
</dependency>
案例一:简单的数值统计(最大值、最小值、平均值、总和)
场景: 统计所有订单的金额情况
import org.elasticsearch.action.search.SearchRequest;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.index.query.QueryBuilders;
import org.elasticsearch.search.aggregations.AggregationBuilders;
import org.elasticsearch.search.aggregations.metrics.*;
import org.elasticsearch.search.builder.SearchSourceBuilder;
public class MetricsAggregationExample {
public static void main(String[] args) {
// 创建搜索请求
SearchRequest searchRequest = new SearchRequest("orders");
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
// 可以添加查询过滤条件(可选)
sourceBuilder.query(QueryBuilders.rangeQuery("create_time")
.gte("2023-01-01")
.lte("2023-12-31"));
// 添加各种聚合
sourceBuilder.aggregation(
AggregationBuilders.avg("avg_amount").field("amount")
);
sourceBuilder.aggregation(
AggregationBuilders.max("max_amount").field("amount")
);
sourceBuilder.aggregation(
AggregationBuilders.min("min_amount").field("amount")
);
sourceBuilder.aggregation(
AggregationBuilders.sum("total_amount").field("amount")
);
sourceBuilder.aggregation(
AggregationBuilders.count("count_orders").field("order_id")
);
sourceBuilder.size(0); // 不返回文档,只返回聚合结果
searchRequest.source(sourceBuilder);
try {
// 执行查询
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
// 解析聚合结果
ParsedAvg avgAmount = response.getAggregations().get("avg_amount");
ParsedMax maxAmount = response.getAggregations().get("max_amount");
ParsedMin minAmount = response.getAggregations().get("min_amount");
ParsedSum totalAmount = response.getAggregations().get("total_amount");
ParsedValueCount countOrders = response.getAggregations().get("count_orders");
System.out.println("平均金额: " + avgAmount.getValue());
System.out.println("最大金额: " + maxAmount.getValue());
System.out.println("最小金额: " + minAmount.getValue());
System.out.println("总金额: " + totalAmount.getValue());
System.out.println("订单总数: " + countOrders.getValue());
} catch (Exception e) {
e.printStackTrace();
}
}
}
案例二:分组统计(Terms聚合)
场景: 按商品分类统计销售额和订单数量
public class TermsAggregationExample {
public static void main(String[] args) {
SearchRequest searchRequest = new SearchRequest("orders");
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
// 按商品分类聚合,并统计每个分类的销售总额和订单数量
TermsAggregationBuilder categoryAgg = AggregationBuilders.terms("category_stats")
.field("category.keyword")
.size(10) // 返回前10个分类
.order(Terms.Order.aggregation("total_amount", false)); // 按总金额降序
// 在每个分类下添加子聚合
categoryAgg.subAggregation(
AggregationBuilders.sum("total_amount").field("amount")
);
categoryAgg.subAggregation(
AggregationBuilders.avg("avg_amount").field("amount")
);
categoryAgg.subAggregation(
AggregationBuilders.cardinality("unique_users").field("user_id")
);
sourceBuilder.aggregation(categoryAgg);
sourceBuilder.size(0);
searchRequest.source(sourceBuilder);
try {
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
Terms terms = response.getAggregations().get("category_stats");
for (Terms.Bucket bucket : terms.getBuckets()) {
String category = bucket.getKeyAsString();
long docCount = bucket.getDocCount();
Sum totalAmount = bucket.getAggregations().get("total_amount");
Avg avgAmount = bucket.getAggregations().get("avg_amount");
Cardinality uniqueUsers = bucket.getAggregations().get("unique_users");
System.out.println("分类: " + category);
System.out.println(" 订单数量: " + docCount);
System.out.println(" 总销售额: " + totalAmount.getValue());
System.out.println(" 平均金额: " + avgAmount.getValue());
System.out.println(" 独立用户数: " + uniqueUsers.getValue());
}
} catch (Exception e) {
e.printStackTrace();
}
}
}
案例三:日期直方图聚合(按时间统计)
场景: 统计每天的订单数量和销售额
public class DateHistogramExample {
public static void main(String[] args) {
SearchRequest searchRequest = new SearchRequest("orders");
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
// 按天聚合统计
DateHistogramAggregationBuilder dateAgg = AggregationBuilders.dateHistogram("daily_stats")
.field("order_date")
.calendarInterval(DateHistogramInterval.DAY) // 按天
.format("yyyy-MM-dd")
.minDocCount(1); // 只显示有数据的天数
// 添加子聚合
dateAgg.subAggregation(
AggregationBuilders.sum("daily_amount").field("amount")
);
dateAgg.subAggregation(
AggregationBuilders.cardinality("daily_users").field("user_id")
);
sourceBuilder.aggregation(dateAgg);
sourceBuilder.size(0);
searchRequest.source(sourceBuilder);
try {
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
ParsedDateHistogram dateHistogram = response.getAggregations().get("daily_stats");
for (Histogram.Bucket bucket : dateHistogram.getBuckets()) {
String date = bucket.getKeyAsString(); // 日期
long docCount = bucket.getDocCount(); // 订单数量
Sum dailyAmount = bucket.getAggregations().get("daily_amount");
Cardinality dailyUsers = bucket.getAggregations().get("daily_users");
System.out.println("日期: " + date);
System.out.println(" 订单数: " + docCount);
System.out.println(" 销售额: " + dailyAmount.getValue());
System.out.println(" 用户数: " + dailyUsers.getValue());
}
} catch (Exception e) {
e.printStackTrace();
}
}
}
案例四:嵌套聚合(多维度分析)
场景: 按年份-季度-月份分层统计,并查看每个季度的TOP商品
public class NestedAggregationExample {
public static void main(String[] args) {
SearchRequest searchRequest = new SearchRequest("orders");
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
// 第一层:按年聚合
DateHistogramAggregationBuilder yearAgg = AggregationBuilders.dateHistogram("year_stats")
.field("order_date")
.calendarInterval(DateHistogramInterval.YEAR)
.format("yyyy");
// 第二层:按季度聚合(嵌套在年份内)
DateHistogramAggregationBuilder quarterAgg = AggregationBuilders.dateHistogram("quarter_stats")
.field("order_date")
.calendarInterval(DateHistogramInterval.QUARTER)
.format("yyyy-MM");
// 第三层:按商品分类聚合
TermsAggregationBuilder productAgg = AggregationBuilders.terms("product_stats")
.field("product_name.keyword")
.size(5) // 每季度TOP5商品
.order(Terms.Order.aggregation("product_amount", false));
productAgg.subAggregation(
AggregationBuilders.sum("product_amount").field("amount")
);
// 构建嵌套结构
quarterAgg.subAggregation(productAgg);
yearAgg.subAggregation(quarterAgg);
sourceBuilder.aggregation(yearAgg);
sourceBuilder.size(0);
// 执行查询和解析(结果解析略,逻辑类似上述案例)
// ...
}
}
案例五:百分比统计(Percentiles)
场景: 分析订单金额的分布情况
public class PercentilesExample {
public static void main(String[] args) {
SearchRequest searchRequest = new SearchRequest("orders");
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
ParseField customRange = new ParseField("percents");
// 计算百分位数
PercentilesAggregationBuilder percentilesAgg =
AggregationBuilders.percentiles("amount_percentiles")
.field("amount")
.percentiles(25, 50, 75, 90, 95, 99);
sourceBuilder.aggregation(percentilesAgg);
sourceBuilder.size(0);
try {
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
ParsedPercentiles percentiles =
response.getAggregations().get("amount_percentiles");
System.out.println("25% 订单金额小于: " + percentiles.percentile(25));
System.out.println("中位数(50%): " + percentiles.percentile(50));
System.out.println("75% 订单金额小于: " + percentiles.percentile(75));
System.out.println("90% 订单金额小于: " + percentiles.percentile(90));
System.out.println("95% 订单金额小于: " + percentiles.percentile(95));
System.out.println("99% 订单金额小于: " + percentiles.percentile(99));
} catch (Exception e) {
e.printStackTrace();
}
}
}
实用建议
- 性能优化:尽量使用
sourceBuilder.size(0)避免返回大量文档数据 - 索引优化:对需要聚合的字段(特别是
keyword类型)设置doc_values: true - 内存控制:使用
size参数限制聚合桶的数量,避免内存溢出 - 查询过滤:先通过
query过滤数据范围,再在有限数据集上进行聚合 - 批量统计:可以在一次请求中组合多个聚合,减少网络请求
这些案例覆盖了常见的统计需求:数值聚合、分组聚合、时间序列分析、嵌套分析和分布分析,根据实际业务需求组合使用即可。