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

如何在PySpark中捕获记录变化?双DataFrame对比实现

PySpark DataFrame 匹配标记解决方案

需求说明

  • 找出df2中同一EmpID、同一日期下Work类型未在df1出现的记录,标记Type为00
  • 找出两DataFrame中同一EmpID、同一日期下的共同Work类型记录:df1的记录标记Type为10,df2的标记为11

数据构建代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import to_date

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

# 构建df1数据
data1 = [
    ("E001", "A", "2024-05-01 09:00:00", "2024-05-01 12:00:00"),
    ("E001", "B", "2024-05-01 13:00:00", "2024-05-01 18:00:00"),
    ("E002", "A", "2024-05-01 08:00:00", "2024-05-01 17:00:00")
]
df1 = spark.createDataFrame(data1, ["EmpID", "Work", "StartTime", "StopTime"])
# 提取日期字段作为匹配维度
df1 = df1.withColumn("WorkDate", to_date("StartTime"))

# 构建df2数据
data2 = [
    ("E001", "A", "2024-05-01 09:30:00", "2024-05-01 12:30:00"),
    ("E001", "C", "2024-05-01 14:00:00", "2024-05-01 19:00:00"),
    ("E002", "A", "2024-05-01 08:30:00", "2024-05-01 17:30:00"),
    ("E002", "B", "2024-05-01 18:00:00", "2024-05-01 20:00:00")
]
df2 = spark.createDataFrame(data2, ["EmpID", "Work", "StartTime", "StopTime"])
df2 = df2.withColumn("WorkDate", to_date("StartTime"))

实现逻辑与代码

核心思路

  1. 从StartTime提取日期字段,作为EmpID、Work之外的第三个匹配维度
  2. 生成df1的唯一匹配键集合(EmpID、WorkDate、Work),用于快速判断Work类型是否存在
  3. 分别处理两类需求:共同记录标记、df2独有记录标记
  4. 合并结果并排序输出

完整实现代码

from pyspark.sql.functions import col, lit, when
from pyspark.sql.types import StringType

# 生成df1的匹配键视图,用于后续判断
df1_keys = df1.select("EmpID", "WorkDate", "Work").distinct().withColumnRenamed("Work", "df1_Work")

# 处理df1的共同记录:标记为10
df1_marked = df1.join(df1_keys, on=["EmpID", "WorkDate", "Work"], how="inner") \
                .withColumn("Type", lit("10").cast(StringType()))

# 处理df2的记录:区分共同记录(标记11)和独有记录(标记00)
df2_joined = df2.join(df1_keys, on=["EmpID", "WorkDate"], how="left")
df2_marked = df2_joined.withColumn(
    "Type",
    when(col("df1_Work") == col("Work"), lit("11"))
    .otherwise(lit("00"))
).drop("df1_Work")

# 合并两个标记后的DataFrame
final_result = df1_marked.unionByName(df2_marked)

# 按指定顺序展示结果
final_result.orderBy("EmpID", "WorkDate", "Work").show(truncate=False)

预期输出示例

+-----+----+-------------------+-------------------+----------+----+
|EmpID|Work|StartTime          |StopTime           |WorkDate  |Type|
+-----+----+-------------------+-------------------+----------+----+
|E001 |A   |2024-05-01 09:00:00|2024-05-01 12:00:00|2024-05-01|10  |
|E001 |A   |2024-05-01 09:30:00|2024-05-01 12:30:00|2024-05-01|11  |
|E001 |B   |2024-05-01 13:00:00|2024-05-01 18:00:00|2024-05-01|10  |
|E001 |C   |2024-05-01 14:00:00|2024-05-01 19:00:00|2024-05-01|00  |
|E002 |A   |2024-05-01 08:00:00|2024-05-01 17:00:00|2024-05-01|10  |
|E002 |A   |2024-05-01 08:30:00|2024-05-01 17:30:00|2024-05-01|11  |
|E002 |B   |2024-05-01 18:00:00|2024-05-01 20:00:00|2024-05-01|00  |
+-----+----+-------------------+-------------------+----------+----+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 02:55:21