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

如何对值不同的PySpark DataFrame进行交并集处理得到指定零值数据集

PySpark 零值数据集合并实现方案

核心逻辑

  1. 从总数据集中筛选出所有value = 0的条目,作为待补充的零值集合
  2. 过滤异常输出表,移除其中在总表中value ≠ 0的对应key
  3. 合并两部分处理后的数据,去重后得到最终结果

完整代码

# 导入依赖
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when

# 初始化SparkSession(如果已有可跳过)
spark = SparkSession.builder.appName("zero_dataset_merge").getOrCreate()

# ---------------------- 步骤1:构造测试数据(实际使用时替换为你的数据源读取代码) ----------------------
# 总数据集
total_data = [("a", 0.5), ("b", 0.4), ("c", 0.5), ("d", 0.3), ("x", 0.0), ("y", 0.0), ("z", 0.0)]
df_total = spark.createDataFrame(total_data, schema=["key", "value"])
# 异常输出数据集
abnormal_data = [("a", 0.0), ("e", 0.0), ("f", 0.0), ("g", 0.0)]
df_abnormal = spark.createDataFrame(abnormal_data, schema=["key", "value"])

# ---------------------- 步骤2:数据处理 ----------------------
# 筛选总表中value为0的条目
df_total_zero = df_total.filter(col("value") == 0)

# 给总表加非零标记,用于过滤异常表
df_total_flag = df_total.withColumn("is_non_zero", when(col("value") != 0, 1).otherwise(0))

# 异常表左连总表标记,过滤掉总表中value非零的key
df_abnormal_filtered = df_abnormal.join(
    df_total_flag.select("key", "is_non_zero"),
    on="key",
    how="left"
).filter(
    (col("is_non_zero").isNull()) | (col("is_non_zero") == 0)
).select("key", "value")

# 合并两部分数据,去重后按key排序
df_final = df_abnormal_filtered.unionByName(df_total_zero).dropDuplicates(["key"]).orderBy("key")

# 输出结果
df_final.show()

输出验证

运行代码后输出结果和预期一致:

+---+-----+
|key|value|
+---+-----+
|  e|  0.0|
|  f|  0.0|
|  g|  0.0|
|  x|  0.0|
|  y|  0.0|
|  z|  0.0|
+---+-----+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 02:18:00