Spark DataFrame按自定义分区(红/非红)计算值占比问题
Got it, let's tackle this problem! You need to split your DataFrame into two custom partitions: one where type equals "red", and another where type is anything else. Then calculate each value's proportion relative to the total of its partition. Here's a straightforward way to do this using Spark's window functions with a derived grouping key:
Step 1: Create a Custom Grouping Column
First, we'll add a new column that acts as our partition key. We'll label rows with type = red as "red_group" and all others as "non_red_group". This gives us the two partitions we need, bypassing the limitation of only partitioning by existing full columns.
Step 2: Use Window Functions to Calculate Group Totals
Next, we'll define a window partitioned by our new grouping column. We'll use the sum() function over this window to get the total value for each group.
Step 3: Compute the Proportion
Finally, divide each row's value by its group's total to get the portion.
Full Scala Code Example
import org.apache.spark.sql.functions.{col, when, sum} import org.apache.spark.sql.expressions.Window // Sample DataFrame val df = spark.createDataFrame(Seq( (1, "red", 10), (2, "blue", 20), (3, "yellow", 30) )).toDF("ID", "type", "value") // Define the custom grouping column val groupedDf = df.withColumn( "group_key", when(col("type") === "red", "red_group").otherwise("non_red_group") ) // Define window partitioned by the custom group key val groupWindow = Window.partitionBy("group_key") // Calculate group total and then the portion val resultDf = groupedDf.withColumn( "total_value", sum(col("value")).over(groupWindow) ).withColumn( "portion", col("value") / col("total_value") ).drop("group_key", "total_value") // Clean up intermediate columns // Show the result resultDf.show()
Full PySpark Code Example
If you're working with Python, the logic is identical—just adjust the syntax:
from pyspark.sql import functions as F from pyspark.sql.window import Window # Sample DataFrame df = spark.createDataFrame([ (1, "red", 10), (2, "blue", 20), (3, "yellow", 30) ], ["ID", "type", "value"]) # Custom grouping column grouped_df = df.withColumn( "group_key", F.when(F.col("type") == "red", "red_group").otherwise("non_red_group") ) # Window definition group_window = Window.partitionBy("group_key") # Calculate portion result_df = grouped_df.withColumn( "total_value", F.sum(F.col("value")).over(group_window) ).withColumn( "portion", F.col("value") / F.col("total_value") ).drop("group_key", "total_value") result_df.show()
Expected Output
+---+------+-----+-------+ | ID| type|value|portion| +---+------+-----+-------+ | 1| red| 10| 1.0| | 2| blue| 20| 0.4| | 3|yellow| 30| 0.6| +---+------+-----+-------+
Quick Explanation
- The
group_keycolumn creates the exact two partitions we need: one for all "red" rows, and one for everything else. - The window function
sum(value).over(groupWindow)calculates the total value for each group. For "red_group", the total is 10, so 10/10 = 1.0. For "non_red_group", total is 50, so 20/50=0.4 and 30/50=0.6. - We drop the intermediate columns to match your desired output structure.
内容的提问来源于stack exchange,提问作者cozyss

