spark: how to groupby a dataframe and to transform each group with

Viewed 126

I have a dataFrame with this columns (site_id,meter_id,timestamp,energy_type).

I'd like grouping by 2 columns (timestamp,energy_type).

Once done, I need to transform every group using a function.

 df.groupby(timestamp,energy_type).<transform_every_group_with_a_function>()

From the groupby, I receive back a RelationalGroupedDataset, how can I transform every group using a function?

thank you

1 Answers

Here's how to pass all groups to a function. If you need help with UDF's here's a simple exmaple.. Generally speaking if you are using a UDF you need to wrap all the columns in a struct to be able to use it. You can access the parameter just like it's a table with columns. You could also just call a map on the output and it would be the same idea, you have access to all the columns to do work on them. That might be better as UDFs don't perform all that great. But the principle is the same, use collect_list with a struct to keep the rows together so you can do work on them.

import spark.implicits._

def convertCase ([list of parameters]) : [return type]

val convertUDF = udf(convertCase)
spark.udf.register("convertUDF", convertCase)

val columns = dataset1.select( dataset1.columns.map(c => col(c)) // create an array of columns
df.groupby(timestamp,energy_type)
.agg(
 collect_list( 
  struct( 
   columns:_* // us a 'splat' to pass all the arguments into a struct
  )
 ).as("all_columns") 
)
.select( 
 myUDF(
  col("all_columns") // pass your struct to a UDF function
 ) 
)
Related