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

PySpark中按条件关联数据集:避免扩行、仅保留匹配首行及添加标记

PySpark:关联数据集时避免扩行并添加存在性标记(仅匹配首行)

需求梳理

核心目标:

  • 给df1每一行添加标记,判断df2中是否存在满足df2.id == df1.id且df2.date >= df1.date的行
  • 若存在匹配,仅关联df2中符合条件的首行(可指定排序逻辑),绝对避免关联后df1的行被扩成多行
  • 禁止先全关联再用窗口函数去重的方案

示例数据准备

先创建测试数据集:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when, first

spark = SparkSession.builder.appName("MatchFirstRow").getOrCreate()

# 构建df1
data1 = [
    ("abc", "London", "football", "2022-02-11"),
    ("def", "Paris", "volley", "2022-02-10"),
    ("ghi", "Manchester", "basketball", "2022-02-09")
]
df1 = spark.createDataFrame(data1, ["id", "city", "sport_event", "date"])
df1 = df1.withColumn("date", col("date").cast("date"))

# 构建df2
data2 = [
    ("abc", 100000, "2022-01-10"),
    ("abc", 200000, "2022-04-15"),
    ("abc", 150000, "2022-02-11")
]
df2 = spark.createDataFrame(data2, ["id", "num_spect", "date"])
df2 = df2.withColumn("date", col("date").cast("date"))

解决方案

方法:左连接+分组聚合(避免扩行)

通过left join匹配符合条件的行,再以df1的所有字段为分组键,用first()聚合df2的字段,直接得到每个df1行对应的唯一匹配项,同时生成存在性标记:

# 若需要指定匹配行的顺序(比如按date升序取第一个),先对df2排序
df2_sorted = df2.orderBy("id", "date")

# 关联+聚合,得到最终结果
result = df1.join(df2_sorted, 
                  (df1.id == df2_sorted.id) & (df2_sorted.date >= df1.date), 
                  how="left") \
            .groupBy(df1.id, df1.city, df1.sport_event, df1.date) \
            .agg(
                first(df2_sorted.num_spect).alias("matched_num_spect"),
                when(first(df2_sorted.id).isNotNull(), 1).otherwise(0).alias("has_match")
            )

result.show()

期望输出

+---+----------+-----------+----------+----------------+----------+
| id|      city|sport_event|      date|matched_num_spect|has_match|
+---+----------+-----------+----------+----------------+----------+
|abc|    London|   football|2022-02-11|           150000|         1|
|def|     Paris|     volley|2022-02-10|            null|         0|
|ghi|Manchester| basketball|2022-02-09|            null|         0|
+---+----------+-----------+----------+----------------+----------+

关键说明

  1. 避免扩行:通过groupBy(df1的所有字段),确保每个分组对应df1的唯一一行,first()只会取该分组内的第一个匹配项
  2. 存在性标记:利用first(df2_sorted.id)是否为null判断是否有匹配,生成has_match标记
  3. 自定义匹配顺序:提前对df2按id和目标字段(比如date)排序,first()就会取排序后的第一个符合条件的行,可根据需求调整排序规则(比如降序取最近的日期)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.15 22:30:54