使用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
相关产品推荐
相关产品推荐

