如何让Elasticsearch而非PySpark执行聚合操作?
如何让Elasticsearch而非PySpark执行聚合操作?
我可以通过以下代码查询Elasticsearch并对结果进行聚合:
spark.read.format("org.elasticsearch.spark.sql").load('my_index').selectExpr('max(my_field)').show()
但这种方式效率很低,因为PySpark会拉取所有数据后再执行聚合。我希望让Elasticsearch完成聚合操作,仅将结果返回给PySpark,该如何实现?
我知道可以通过以下方式让Elasticsearch执行过滤,仅返回过滤后的结果:
myquery = '''{ "range": { "my_field": { "gt": 1, "boost": 1 } } }''' spark.read.format("org.elasticsearch.spark.sql").option("es.query", myquery).load('my_index').show()
有趣的是,当我使用PySpark的filter语句时:
spark.read.format("org.elasticsearch.spark.sql").load('my_index').filter('my_field>1').show()
似乎也是由Elasticsearch执行过滤。但我始终无法实现让Elasticsearch执行聚合操作。
解决方案
方法1:手动编写带聚合的ES DSL查询
直接在Elasticsearch的查询DSL中定义聚合逻辑,让ES提前计算完成后再返回结果给Spark,这是最可控的方式:
- 编写包含聚合的DSL(以求
my_field最大值为例):
{ "size": 0, // 不返回原始文档,仅返回聚合结果 "aggs": { "max_my_field": { "max": { "field": "my_field" } } } }
- 在PySpark中传入该查询,并指定读取聚合结果的字段路径:
agg_query = '''{ "size": 0, "aggs": { "max_my_field": { "max": { "field": "my_field" } } } }''' df = spark.read.format("org.elasticsearch.spark.sql")\ .option("es.query", agg_query)\ .option("es.read.field.as.array.include", "aggregations.max_my_field.value")\ .load("my_index") df.show()
这种方式下,ES只会返回聚合后的单个数值,完全避免拉取全量数据。
方法2:开启聚合下推优化(版本依赖)
从elasticsearch-spark-sql 7.x及以上版本开始,Spark的部分聚合操作会自动下推到ES执行,无需手动编写DSL。你需要先确认版本匹配,然后开启相关配置:
# 开启Spark和ES的聚合下推 spark.conf.set("spark.sql.pushdown.aggregate", "true") spark.conf.set("spark.sql.es.pushdown.aggregate", "true") # 原聚合代码即可自动下推到ES执行 spark.read.format("org.elasticsearch.spark.sql").load('my_index').selectExpr('max(my_field)').show()
验证聚合是否在ES端执行
可以通过两种方式确认:
- 查看Elasticsearch集群的查询日志,检查是否有聚合请求
- 查看Spark UI的Stage详情,如果输入数据量远小于索引总数据量,说明是ES聚合后的结果
内容的提问来源于stack exchange,提问作者Arik
相关产品推荐
相关产品推荐

