如何用PySpark合并DataFrame中Flag1与Flag2列并按规则聚合?
需求说明
按id、status、type分组合并行,根据给定的Flag推导规则填充Flag1和Flag2的null值,最终确保两者不同时为空。
实现步骤
- 创建测试DataFrame(已有数据可跳过此步)
from pyspark.sql import SparkSession from pyspark.sql import functions as F spark = SparkSession.builder.appName("MergeFlagRecords").getOrCreate() # 模拟原始数据 data = [ (123, "N/A", None, False, "Prepaid"), (123, "N/A", True, None, "Prepaid") ] df = spark.createDataFrame(data, schema=["id", "status", "Flag1", "Flag2", "type"])
- 分组聚合收集非null Flag值
按id、status、type分组,提取每组内非null的Flag1和Flag2唯一有效值:
grouped_df = df.groupBy("id", "status", "type") \ .agg( F.collect_set(F.when(F.col("Flag1").isNotNull(), F.col("Flag1"))).alias("Flag1_vals"), F.collect_set(F.when(F.col("Flag2").isNotNull(), F.col("Flag2"))).alias("Flag2_vals") ) \ .withColumn("Flag1", F.element_at("Flag1_vals", 1)) \ .withColumn("Flag2", F.element_at("Flag2_vals", 1))
- 按规则填充null值
利用when函数匹配你给出的推导规则,填充其中一个Flag为空的情况,同时处理反向逻辑:
final_df = grouped_df \ # 反向推导:Flag1为空时,根据Flag2的值填充 .withColumn( "Flag1", F.when(F.col("Flag1").isNull(), F.when(F.col("Flag2") == True, True).when(F.col("Flag2") == False, False) ).otherwise(F.col("Flag1")) ) \ # 正向推导:Flag2为空时,根据Flag1的值填充 .withColumn( "Flag2", F.when(F.col("Flag2").isNull(), F.when(F.col("Flag1") == True, True).when(F.col("Flag1") == False, False) ).otherwise(F.col("Flag2")) ) \ .drop("Flag1_vals", "Flag2_vals") \ # 过滤掉Flag1和Flag2同时为空的行 .filter(F.col("Flag1").isNotNull() | F.col("Flag2").isNotNull())
- 查看最终结果
final_df.show()
输出结果:
+---+------+-------+-----+-----+ | id|status| type|Flag1|Flag2| +---+------+-------+-----+-----+ |123| N/A|Prepaid| true|false| +---+------+-------+-----+-----+
规则匹配说明
完全贴合你给出的规则逻辑:
- 正向推导(Flag1→Flag2):
True→null → 填充Flag2为True;False→null → 填充Flag2为False;非空值直接保留 - 反向推导(Flag2→Flag1):
null→True → 填充Flag1为True;null→False → 填充Flag1为False;非空值直接保留
内容的提问来源于stack exchange,提问作者veganzombie
相关产品推荐
相关产品推荐

