如何对值不同的PySpark DataFrame进行交并集处理得到指定零值数据集
PySpark 零值数据集合并实现方案
核心逻辑
- 从总数据集中筛选出所有
value = 0的条目,作为待补充的零值集合 - 过滤异常输出表,移除其中在总表中
value ≠ 0的对应key - 合并两部分处理后的数据,去重后得到最终结果
完整代码
# 导入依赖 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
相关产品推荐
相关产品推荐

