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

基于其他列聚合创建新列——PySpark实现方案

PySpark实现分组提取指定条件值并填充列

你可以通过以下两种方法实现需求:

方法一:使用窗口函数(推荐)

利用窗口函数在分组内提取share=0.01对应的units值,并将该值填充到整个分组的first列:

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 定义窗口:按gender、week、type分组
window_spec = Window.partitionBy("gender", "week", "type")

# 添加first列
result_df = df.withColumn(
    "first",
    F.first(F.when(F.col("share") == 0.01, F.col("units")), ignorenulls=True).over(window_spec)
)

# 查看结果
result_df.show()

逻辑说明:

  • F.when(F.col("share") == 0.01, F.col("units")):仅当share等于0.01时保留units值,其余情况返回null
  • F.first(..., ignorenulls=True):在分组内忽略null值,取第一个非null的units值(即share=0.01对应的数值)
  • 窗口函数自动将该值广播到分组内所有行,完成整组填充

方法二:分组筛选后关联

先筛选出share=0.01的行,再通过关联将结果合并到原DataFrame:

from pyspark.sql import functions as F

# 提取每个分组内share=0.01的units值,重命名为first
grouped_df = df.filter(F.col("share") == 0.01).select(
    "gender", "week", "type",
    F.col("units").alias("first")
)

# 左关联原DataFrame,填充first列
result_df = df.join(grouped_df, on=["gender", "week", "type"], how="left")

# 查看结果
result_df.show()

逻辑说明:

  • 第一步筛选出所有share=0.01的行,每个分组仅保留一行数据
  • 左关联操作将first列的值匹配到原DataFrame对应分组的所有行中

内容的提问来源于stack exchange,提问作者paulo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 16:40:22