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

PySpark数据处理:用同币种非0值填充列并新增列

解决PySpark DataFrame替换同组0值问题

要实现将RATIO_FROM和RATIO_TO中的0值替换为对应FROM_CURRENCY分组下的非0值,你需要用窗口函数来动态获取同组的非0基准值——F.lit()只能传入固定常量,无法根据每行的分组匹配动态数据,这就是你之前尝试失败的原因。

具体实现步骤

  1. 导入所需的PySpark函数
  2. 定义窗口规则:按FROM_CURRENCY分区,确保非0值优先被选取
  3. 使用when()函数判断原列是否为0,替换为窗口计算的基准值,否则保留原值

完整代码示例

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

# 初始化SparkSession(如果未初始化)
spark = SparkSession.builder.appName("ReplaceZeroValues").getOrCreate()

# 模拟你的原始DataFrame
data = [
    ("AED", "EUR", 0, 0),
    ("AED", "EUR", 1, 1),
    ("GNF", "EUR", 0, 0),
    ("DZD", "EUR", 1, 1),
    ("GNF", "EUR", 1000, 1000)
]
df = spark.createDataFrame(data, ["FROM_CURRENCY", "TO_CURRENCY", "RATIO_FROM", "RATIO_TO"])

# 定义窗口:按FROM_CURRENCY分区,让非0值排在最前面
window_spec = Window.partitionBy("FROM_CURRENCY").orderBy(F.when(F.col("RATIO_FROM") != 0, 1).desc())

# 生成目标新列
df_result = df.withColumn(
    "RATIO_FROM_BIS",
    F.when(F.col("RATIO_FROM") == 0, F.first("RATIO_FROM").over(window_spec)).otherwise(F.col("RATIO_FROM"))
).withColumn(
    "RATIO_TO_BIS",
    F.when(F.col("RATIO_TO") == 0, F.first("RATIO_TO").over(window_spec)).otherwise(F.col("RATIO_TO"))
)

# 查看最终结果
df_result.select("FROM_CURRENCY", "TO_CURRENCY", "RATIO_FROM_BIS", "RATIO_TO_BIS").show()

代码说明

  • 窗口规则:orderBy(F.when(F.col("RATIO_FROM") != 0, 1).desc())会把同组内非0值的行排在最前面,first()就能直接取到该组的基准非0值。
  • when()逻辑:判断原列是否为0,是则替换为窗口获取的基准值,否则保留原值。
  • 若同组内存在多个不同的非0值,可根据业务需求替换first()为max()、min()或avg()等聚合函数。

替代方案:分组聚合后Join

如果窗口函数理解起来有难度,也可以先分组计算每个FROM_CURRENCY的非0基准值,再和原表关联:

# 分组提取每个FROM_CURRENCY的非0基准值
base_values = df.filter(F.col("RATIO_FROM") != 0).select(
    "FROM_CURRENCY",
    F.first("RATIO_FROM").alias("BASE_FROM"),
    F.first("RATIO_TO").alias("BASE_TO")
).distinct()

# 关联后替换0值
df_result = df.join(base_values, on="FROM_CURRENCY", how="left").withColumn(
    "RATIO_FROM_BIS",
    F.when(F.col("RATIO_FROM") == 0, F.col("BASE_FROM")).otherwise(F.col("RATIO_FROM"))
).withColumn(
    "RATIO_TO_BIS",
    F.when(F.col("RATIO_TO") == 0, F.col("BASE_TO")).otherwise(F.col("RATIO_TO"))
).drop("BASE_FROM", "BASE_TO")

内容的提问来源于stack exchange,提问作者f.ivy

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 10:06:22