Spark Scala 3.1按指定后缀分组列生成嵌套Struct DataFrame
Answer
Here's a clean, dynamic implementation in Scala for Spark 3.1 that transforms your DataFrame exactly as requested:
import org.apache.spark.sql.functions.{col, struct} import org.apache.spark.sql.DataFrame def transformToNestedTags(df: DataFrame, inputSuffix: Array[String]): DataFrame = { // Create a struct for each suffix group val suffixStructs = inputSuffix.map { suffix => // Collect key columns (key*_suffix) first, then add the standalone suffix column val keyColumns = df.columns.filter(c => c != suffix && c.endsWith(suffix)).map(col) val allGroupColumns = keyColumns :+ col(suffix) // Wrap into a struct and alias with the suffix name struct(allGroupColumns: _*).alias(suffix) } // Combine all suffix structs into the top-level "tags" struct val tagsColumn = struct(suffixStructs: _*).alias("tags") // Select only the id and new tags column df.select(col("id"), tagsColumn) } // Example usage: val inputSuffix = Array("suffix1", "suffix2") val originalDF = spark.read... // Add your DataFrame loading logic here val transformedDF = transformToNestedTags(originalDF, inputSuffix) // Verify the resulting schema transformedDF.printSchema()
How this works:
- Dynamic column grouping: For each suffix in your input array, we filter the DataFrame's columns to capture all related fields (key*_suffix plus the standalone suffix column). We explicitly order the key columns first followed by the suffix column to match your desired schema structure.
- Nested struct creation: Spark's
structfunction wraps each group of columns into a nested struct, aliased with the suffix name. - Top-level tags struct: All suffix-specific structs are combined into a single
tagscolumn, creating the nested hierarchy you need. - Reusable function: This implementation works for any number of suffixes, not just the two in your example.
When you run transformedDF.printSchema(), you'll see the exact schema you specified in your question.
内容的提问来源于stack exchange,提问作者Yasha Jadwani
相关产品推荐
相关产品推荐

