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

PySpark合并DataFrame:优先保留月度结算值,无则显示日估算值

PySpark解决方案:筛选无对应月度结算的日估算值

核心思路

  1. 为两个DataFrame提取年-月维度标识,用于匹配ID对应的月份是否存在结算记录
  2. 标记df2中哪些记录的(ID, 年-月)组合在df1中不存在
  3. 过滤出符合条件的df2记录,并更新DSC字段的说明文本
  4. 合并df1全量数据和过滤后的df2数据,得到最终结果

具体实现代码

首先导入必要函数并创建示例DataFrame(若已有数据可跳过此部分):

from pyspark.sql import SparkSession
from pyspark.sql.functions import date_format, col, when, exists, lit

# 初始化SparkSession
spark = SparkSession.builder.appName("monthly_vs_daily").getOrCreate()

# 创建月度结算值表df1
data1 = [
    ("2022-01-31", 123, 10, "CLOSED MONTH"),
    ("2022-02-28", 123, 20, "CLOSED MONTH"),
    ("2022-03-31", 999, 30, "CLOSED MONTH"),
    ("2022-04-30", 999, 40, "CLOSED MONTH")
]
df1 = spark.createDataFrame(data1, ["DATA", "ID", "VALUE", "DSC"])

# 创建日估算值表df2
data2 = [
    ("2022-01-31", 123, 50, "ESTIMATED DAY"),
    ("2022-02-28", 123, 60, "ESTIMATED DAY"),
    ("2022-03-31", 123, 70, "ESTIMATED DAY"),
    ("2022-04-30", 123, 80, "ESTIMATED DAY"),
    ("2022-03-20", 123, 90, "ESTIMATED DAY"),
    ("2022-03-25", 123, 100, "ESTIMATED DAY"),
    ("2022-04-30", 999, 120, "ESTIMATED DAY"),
    ("2022-05-02", 999, 150, "ESTIMATED DAY"),
    ("2022-05-03", 999, 200, "ESTIMATED DAY")
]
df2 = spark.createDataFrame(data2, ["DATA", "ID", "VALUE", "DSC"])

接下来处理核心逻辑:

# 1. 为两个DF添加年-月字段,统一月份匹配维度
df1_with_month = df1.withColumn("year_month", date_format(col("DATA"), "yyyy-MM"))
df2_with_month = df2.withColumn("year_month", date_format(col("DATA"), "yyyy-MM"))

# 2. 创建df1的(ID, year_month)临时视图,用于存在性判断
df1_with_month.createOrReplaceTempView("closed_month_records")

# 3. 过滤df2:保留(ID, year_month)不在df1中的记录,同时更新DSC字段
filtered_df2 = df2_with_month.filter(
    ~exists(
        spark.table("closed_month_records"),
        lambda cm: cm.ID == col("ID") and cm.year_month == col("year_month")
    )
).withColumn(
    "DSC",
    when(
        exists(
            spark.table("closed_month_records"),
            lambda cm: cm.year_month == col("year_month")
        ),
        lit("ESTIMATED DAY -Because closed month ") + date_format(col("DATA"), "M") + lit(" has different ID")
    ).otherwise(
        lit("ESTIMATED DAY -Because there is no closed month ") + date_format(col("DATA"), "M")
    )
).drop("year_month")

# 4. 合并df1和过滤后的df2
final_df = df1.unionByName(filtered_df2)

# 查看排序后的结果
final_df.orderBy("ID", "DATA").show(truncate=False)

代码说明

  • 提取年-月:通过date_format将日期转换为yyyy-MM格式,确保月份匹配的一致性
  • 存在性判断:用exists函数检查df2记录的(ID, 年-月)是否在df1中存在,~表示取反,仅保留无对应结算的记录
  • DSC字段更新:分两种场景生成说明:
    • 若该月份有其他ID的结算记录,生成"同月份不同ID"的说明
    • 若该月份无任何结算记录,生成"无对应月度结算"的说明
  • 合并数据:使用unionByName确保列名一致的前提下合并两个DataFrame,最后按ID和日期排序输出

输出结果

运行代码后将得到与预期一致的结果:

+----------+---+-----+---------------------------------------------------+
|DATA      |ID |VALUE|DSC                                               |
+----------+---+-----+---------------------------------------------------+
|2022-01-31|123|10   |CLOSED MONTH                                      |
|2022-02-28|123|20   |CLOSED MONTH                                      |
|2022-03-20|123|90   |ESTIMATED DAY -Because closed month 3 has different ID|
|2022-03-25|123|100  |ESTIMATED DAY -Because closed month 3 has different ID|
|2022-03-31|999|30   |CLOSED MONTH                                      |
|2022-04-30|999|40   |CLOSED MONTH                                      |
|2022-05-02|999|150  |ESTIMATED DAY -Because there is no closed month 5 |
|2022-05-03|999|200  |ESTIMATED DAY -Because there is no closed month 5 |
+----------+---+-----+---------------------------------------------------+

内容的提问来源于stack exchange,提问作者Gustavo Morais Oliveira

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 09:50:17