You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame按自定义分区(红/非红)计算值占比问题

How to Calculate Value Proportions with Custom Partitioning in 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_key column 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 04:31:09