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| +---+----------+-----------+----------+----------------+----------+
关键说明
- 避免扩行:通过
groupBy(df1的所有字段),确保每个分组对应df1的唯一一行,first()只会取该分组内的第一个匹配项 - 存在性标记:利用
first(df2_sorted.id)是否为null判断是否有匹配,生成has_match标记 - 自定义匹配顺序:提前对
df2按id和目标字段(比如date)排序,first()就会取排序后的第一个符合条件的行,可根据需求调整排序规则(比如降序取最近的日期)
内容的提问来源于stack exchange,提问作者Jresearcher
相关产品推荐
相关产品推荐

