I have a dataset with heterogeneous JSONs that I can group by type and apply a StructType to. E.g. RDD[(Type, JSON)] and Set[Type], containing all types from the original RDD.
Now I want to write these JSONs into a typed Parquet files, partitioned by type.
What I do currently is:
val getSchema: Type => StructType = ???
val buildRow: StructType => JSON => Row = ???
types.foreach { jsonType =>
val sparkSchema: StructType = getSchema(jsonType)
val rows: RDD[Row] = rdd
.filter(k => k == jsonType)
.map { case (_, json) => buildRow(sparkSchema)(json) }
spark.createDataFrame(rows, sparkSchema).write.parquet(path)
}
It works, but very inefficiently as it needs to traverse the original RDD many times - I usually have tens and even hundreds of types.
Are there any better solutions? I tried to union the dataframes, but it fails because rows with different arities cannot be coalesced. Or maybe at least I can release the resources taken by a dataframe by unpersisting it?