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

使用Java将Elasticsearch聚合查询结果导出为CSV求助

导出Elasticsearch聚合结果到CSV(Java实现)

问题描述

我正在使用Java查询Elasticsearch,希望将查询的聚合结果导出为CSV文件,恳请提供代码帮助。

原查询代码

try {
    RangeQueryBuilder rangeQ = QueryBuilders
            .rangeQuery("@timestamp")
            .gte("1663632000000")
            .lte("1663804799000")
            .format("epoch_millis");

    TermsAggregationBuilder termsAggregation = AggregationBuilders
            .terms("term_by_client_id")
            .field("labels.client_id")
            .size(100000)
            .minDocCount(1);

    termsAggregation
            .subAggregation(
                    AggregationBuilders
                            .sum("sum_by")
                            .field("labels.row_count")
            );
    termsAggregation
            .subAggregation(
                    AggregationBuilders
                            .terms("term_By_job")
                            .field("labels.job_id")
            );


    SearchRequest searchRequest = new SearchRequest();
    searchRequest.indices("*itm*");


    SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder();
    searchSourceBuilder.query(rangeQ);
    searchSourceBuilder.aggregation(termsAggregation);
    // searchSourceBuilder.size(100000);
    searchRequest.source(searchSourceBuilder);

    SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);
    System.out.println(searchResponse);

    Aggregations aggregations = searchResponse.getAggregations();
    Map<String, Aggregation> aggregationMap = aggregations.asMap();
    for (Map.Entry<String, Aggregation> each : aggregationMap.entrySet()){
        System.out.println((each.getValue()));
    }

} catch (IOException e) {
    throw new RuntimeException(e);
}

聚合结果片段

"buckets":[{"key":"1741433","doc_count":1},{"key":"1741435","doc_count":1},{"key":"1741436","doc_count":1},{"key":"1741440","doc_count":1},{"key":"1741441","doc_count":1},{"key":"1741442","doc_count":1},{"key":"1741443","doc_count":1},{"key":"1741444","doc_count":1},{"key":"1741450","doc_count":1},{"key":"1741451","doc_count":1}]},"sum#sum_by":{"value":1.0951264E7}},{"key":"86206","doc_count":383,"sterms#term_By_job":{"doc_count_error_upper_bound":6,"sum_other_doc_count":361,"buckets":[{"key":"1211310","doc_count":3},{"key":"1211316","doc_count":3},{"key":"1210943","doc_count":2},{"key":"1210945","doc_count":2},{"key":"1210946","doc_count":2},{"key":"1210947","doc_count":2},{"key":"1210948","doc_count":2},{"key":"1210949","doc_count":2},{"key":"1210987","doc_count":2},{"key":"1211010","doc_count":2}]}

解决方案

1. 依赖准备

使用opencsv库简化CSV写入操作,Maven依赖如下:

<dependency>
    <groupId>com.opencsv</groupId>
    <artifactId>opencsv</artifactId>
    <version>5.6</version>
</dependency>

2. 完整导出代码

将原查询代码中的聚合结果循环替换为以下逻辑,实现CSV导出:

try {
    RangeQueryBuilder rangeQ = QueryBuilders
            .rangeQuery("@timestamp")
            .gte("1663632000000")
            .lte("1663804799000")
            .format("epoch_millis");

    TermsAggregationBuilder termsAggregation = AggregationBuilders
            .terms("term_by_client_id")
            .field("labels.client_id")
            .size(100000)
            .minDocCount(1);

    termsAggregation
            .subAggregation(
                    AggregationBuilders
                            .sum("sum_by")
                            .field("labels.row_count")
            );
    termsAggregation
            .subAggregation(
                    AggregationBuilders
                            .terms("term_By_job")
                            .field("labels.job_id")
            );


    SearchRequest searchRequest = new SearchRequest();
    searchRequest.indices("*itm*");

    SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder();
    searchSourceBuilder.query(rangeQ);
    searchSourceBuilder.aggregation(termsAggregation);
    searchRequest.source(searchSourceBuilder);

    SearchResponse searchResponse = client.search(searchRequest, RequestOptions.DEFAULT);

    // -------------------------- 新增CSV导出逻辑 --------------------------
    // 定义CSV文件路径
    String csvFilePath = "./es_aggregation_result.csv";
    CSVWriter writer = new CSVWriter(new FileWriter(csvFilePath));

    // 写入表头
    String[] header = {"client_id", "doc_count", "sum_row_count", "job_id"};
    writer.writeNext(header);

    Aggregations aggregations = searchResponse.getAggregations();
    // 获取client_id的聚合桶
    Terms clientTerms = aggregations.get("term_by_client_id");

    for (Terms.Bucket clientBucket : clientTerms.getBuckets()) {
        String clientId = clientBucket.getKeyAsString();
        long docCount = clientBucket.getDocCount();

        // 获取sum聚合结果
        Sum sumAgg = clientBucket.getAggregations().get("sum_by");
        double sumValue = sumAgg.getValue();

        // 获取job_id的聚合桶
        Terms jobTerms = clientBucket.getAggregations().get("term_By_job");

        // 方式1:每个job_id单独生成一行记录
        for (Terms.Bucket jobBucket : jobTerms.getBuckets()) {
            String jobId = jobBucket.getKeyAsString();
            String[] row = {
                    clientId,
                    String.valueOf(docCount),
                    String.valueOf(sumValue),
                    jobId
            };
            writer.writeNext(row);
        }

        // 方式2:将同一client的所有job_id用逗号分隔,写入同一行
        // List<String> jobIds = new ArrayList<>();
        // for (Terms.Bucket jobBucket : jobTerms.getBuckets()) {
        //     jobIds.add(jobBucket.getKeyAsString());
        // }
        // String jobIdsStr = String.join(",", jobIds);
        // String[] row = {
        //         clientId,
        //         String.valueOf(docCount),
        //         String.valueOf(sumValue),
        //         jobIdsStr
        // };
        // writer.writeNext(row);
    }

    // 关闭写入流
    writer.close();
    System.out.println("CSV文件导出完成,路径:" + csvFilePath);

} catch (IOException e) {
    throw new RuntimeException(e);
}

代码说明

  • 聚合解析:通过Terms和Sum类直接获取Elasticsearch返回的聚合结果,避免手动解析JSON
  • CSV写入:使用CSVWriter自动处理CSV格式(如逗号转义、换行等)
  • 两种写入方式:
    • 方式1:每个job_id对应一行记录,适合需要明细数据的场景
    • 方式2:同一client的job_id拼接成字符串,适合按client汇总的场景

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 19:55:17