Spark Plan for collect_list

Viewed 480

This is with reference Jacek's answer to How to get the size of result generated using concat_ws?

The DSL query in the answer calls collect_list twice with concat and size separately.

input.groupBy($"col1").agg(
     concat_ws(",", collect_list($"COL2".cast("string"))) as "concat",
     size(collect_list($"COL2".cast("string"))) as "size"
)

With anoptimizedPlan like :

Aggregate [COL1#9L], 
[COL1#9L,
concat_ws(,,(hiveudaffunction(HiveFunctionWrapper(GenericUDAFCollectList,GenericUDAFCollectList@470a4e26),cast(COL2#10L as string),false,0,0),mode=Complete,isDistinct=false)) AS concat#13,
size((hiveudaffunction(HiveFunctionWrapper(GenericUDAFCollectList,GenericUDAFCollectList@2602f45),cast(COL2#10L as string),false,0,0),mode=Complete,isDistinct=false)) AS size#14]
+- Project [(id#8L % 2) AS COL1#9L,id#8L AS COL2#10L]
   +- LogicalRDD [id#8L], MapPartitionsRDD[20] at range at <console>:25

How will it be different performance-wise, if I use collect_list only once and then withColumn API to generate the other two columns?

input
  .groupBy("COL1")
  .agg(collect_list($"COL2".cast("string")).as("list") )
  .withColumn("concat", concat_ws("," , $"list"))
  .withColumn("size", size($"list"))
  .drop("list")

Which has an optimizedPlan like :

Project [COL1#9L,
concat_ws(,,list#17) AS concat#18,
size(list#17) AS size#19]
+- Aggregate [COL1#9L], 
[COL1#9L,(hiveudaffunction(HiveFunctionWrapper(GenericUDAFCollectList,GenericUDAFCollectList@5cb88b6b),
   cast(COL2#10L as string),false,0,0),mode=Complete,isDistinct=false) AS list#17]
   +- Project [(id#8L % 2) AS COL1#9L,id#8L AS COL2#10L]
      +- LogicalRDD [id#8L], MapPartitionsRDD[20] at range at <console>:25

I see collect_list being called twice in the former example but just wanted to know if there are any significant differences apart from that. Using Spark 1.6.

0 Answers
Related