Spark Scala实现单列转列需求(Spark2.1.1/Scala2.11.8)
Got it, let's work through this problem together. You want to pivot the values in the C4 column into separate columns, with a 1 marking where the original row matches the column value and 0 otherwise. Here's a step-by-step solution using Spark 2.1.1 and Scala 2.11.8:
First, let's create the sample DataFrame matching your input data:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ // Initialize Spark session val spark = SparkSession.builder() .appName("C4PivotTransformation") .master("local[*]") // Remove this line if running on a distributed cluster .getOrCreate() // Create the original dataset as a DataFrame val originalDF = Seq( ("a", "b", "c", 9), ("x", "y", "z", 6), ("A", "B", "C", 3) ).toDF("C1", "C2", "C3", "C4") // Verify the original data originalDF.show()
This will print your original table:
+---+---+---+---+ |C1 |C2 |C3 |C4 | +---+---+---+---+ |a |b |c |9 | |x |y |z |6 | |A |B |C |3 | +---+---+---+---+
Spark 2.1.1 works best when we explicitly specify the values we want to pivot (this avoids extra scans of the dataset). Let's fetch and sort the unique values from C4:
// Get sorted unique values from C4 to ensure column order matches your example val c4UniqueValues = originalDF.select("C4") .distinct() .collect() .map(_.getInt(0)) .sorted // This gives us [3,6,9] as in your desired output
Now we'll group by the non-pivot columns (C1, C2, C3), pivot on C4, and replace any null values (which mean no match) with 0:
val transformedDF = originalDF.groupBy("C1", "C2", "C3") .pivot("C4", c4UniqueValues) // Use our precomputed values for pivot .agg(count("C4")) // Count returns 1 for matching rows, null otherwise .na.fill(0) // Replace nulls with 0 to get the one-hot pattern // Check the final result transformedDF.show()
This will output exactly the format you requested:
+---+---+---+---+---+---+ |C1 |C2 |C3 |3 |6 |9 | +---+---+---+---+---+---+ |A |B |C |1 |0 |0 | |x |y |z |0 |1 |0 | |a |b |c |0 |0 |1 | +---+---+---+---+---+---+
Quick Notes for Edge Cases:
- If your dataset has duplicate rows for the same
C1/C2/C3combination, swapcount("C4")withmax(lit(1))to ensure you still get a 1 for matches instead of a count of duplicates. - If you don't care about the order of the pivot columns, you can skip the
.sortedstep when fetchingc4UniqueValues.
内容的提问来源于stack exchange,提问作者punati

