Java ES聚合案例怎么统计数据

wen java案例 26

本文目录导读:

Java ES聚合案例怎么统计数据

  1. 环境准备
  2. 案例一:简单的数值统计(最大值、最小值、平均值、总和)
  3. 案例二:分组统计(Terms聚合)
  4. 案例三:日期直方图聚合(按时间统计)
  5. 案例四:嵌套聚合(多维度分析)
  6. 案例五:百分比统计(Percentiles)
  7. 实用建议

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

实用建议

  1. 性能优化:尽量使用sourceBuilder.size(0)避免返回大量文档数据
  2. 索引优化:对需要聚合的字段(特别是keyword类型)设置doc_values: true
  3. 内存控制:使用size参数限制聚合桶的数量,避免内存溢出
  4. 查询过滤:先通过query过滤数据范围,再在有限数据集上进行聚合
  5. 批量统计:可以在一次请求中组合多个聚合,减少网络请求

这些案例覆盖了常见的统计需求:数值聚合、分组聚合、时间序列分析、嵌套分析和分布分析,根据实际业务需求组合使用即可。

抱歉,评论功能暂时关闭!