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

请求修改PySpark代码移除空值(已解决Databricks Merge冲突问题)

解决Merge冲突后移除空值的PySpark代码修改方案

原代码实现了按P_key分区、取Id排序后第一行的去重逻辑,现在添加空值移除逻辑,可根据业务需求选择以下场景处理:

场景1:过滤核心关联键P_key为空的行

空的P_key无合并意义,优先过滤这类行后再执行去重:

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number

# 先过滤P_key为空的行
filtered_df = partdf.filter(partdf.P_key.isNotNull())
# 执行原有去重逻辑
df2 = filtered_df.withColumn("rn", row_number().over(Window.partitionBy("P_key").orderBy("Id")))
df3 = df2.filter("rn = 1").drop("rn")

场景2:同时过滤指定业务字段的空值

如果需要移除其他字段(如Id)的空值,可扩展过滤条件:

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number

# 过滤P_key和Id均不为空的行
filtered_df = partdf.filter(partdf.P_key.isNotNull() & partdf.Id.isNotNull())
df2 = filtered_df.withColumn("rn", row_number().over(Window.partitionBy("P_key").orderBy("Id")))
df3 = df2.filter("rn = 1").drop("rn")

场景3:批量移除所有含空值的行

若业务允许移除任何包含空值的行,可使用dropna()简化操作:

from pyspark.sql.window import Window
from pyspark.sql.functions import row_number

# 移除所有含空值的行
filtered_df = partdf.dropna()
df2 = filtered_df.withColumn("rn", row_number().over(Window.partitionBy("P_key").orderBy("Id")))
df3 = df2.filter("rn = 1").drop("rn")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 13:35:19