In Spark 3.1.1, I did a groupBy without distinct on a DataFrame. I tried to turn off the partial aggregation using
spark.conf.set("spark.sql.aggregate.partialaggregate.skip.enabled", "true")
and then ran the query
df.groupBy("method").agg(sum("request_body_len"))
Spark still ends up doing the partial aggreagtion as shown in the physical plan.
== Physical Plan ==
*(2) HashAggregate(keys=[method#23], functions=[sum(cast(request_body_len#28 as bigint))], output=[method#23, sum(request_body_len)#142L])
+- Exchange hashpartitioning(method#23, 200), ENSURE_REQUIREMENTS, [id=#58]
+- *(1) HashAggregate(keys=[method#23], functions=[partial_sum(cast(request_body_len#28 as bigint))], output=[method#23, sum#146L])
+- FileScan csv [method#23,request_body_len#28] Batched: false, DataFilters: [], Format: CSV, Location: InMemoryFileIndex[file:/mnt/http.csv], PartitionFilters: [], PushedFilters: [], ReadSchema: struct<method:string,request_body_len:int>
I tried this after viewing this video on Youtube : https://youtu.be/_Ne27JcLnEc @56:53
Is this feature no longer available in the latest Spark or am I missing something?