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

如何用Spark Join操作合并PySpark DataFrame得到目标结果?

如何通过Join或高效方式合并两个PySpark DataFrame得到目标结果?

输入DataFrame1

li = [('abc', 'xyz')]
liColumns = ["aid", "bid"]
tempDF = spark.createDataFrame(data=li, schema=liColumns)
tempDF.printSchema()
tempDF.show(truncate=False)

输出:

root
 |-- aid: string (nullable = true)
 |-- bid: string (nullable = true)

+---+---+
|aid|bid|
+---+---+
|abc|xyz|
+---+---+

输入DataFrame2

other_li = [('abc', '111', 'desc111'), ('abc', '112', 'desc112'), ('xyz', 'A123', 'city'), ('xyz', 'A456', 'state'), ('xyz', 'A789', 'zip')]
otherColumns = ['real_aid', 'code', 'some_value']
otherDF = spark.createDataFrame(data=other_li, schema=otherColumns)
otherDF.printSchema()
otherDF.show(truncate=False)

输出:

root
 |-- real_aid: string (nullable = true)
 |-- code: string (nullable = true)
 |-- some_value: string (nullable = true)

+--------+----+----------+
|real_aid|code|some_value|
+--------+----+----------+
|abc     |111 |desc111   |
|abc     |112 |desc112   |
|xyz     |A123|city      |
|xyz     |A456|state     |
|xyz     |A789|zip       |
+--------+----+----------+

问题描述

我需要合并这两个DataFrame得到目标结果,已知可以用append/union实现,但能不能用Join操作完成?有没有更高效的方法?要处理两张大数据表。

预期目标DataFrame

output_li = [('abc', '111', 'desc111'), ('abc', '112', 'desc112'), ('abc', 'A123', 'city'), ('abc', 'A456', 'state'), ('abc', 'A789', 'zip'), ('xyz', 'A123', 'city'), ('xyz', 'A456', 'state'), ('xyz', 'A789', 'zip')]
otherColumns = ['real_aid', 'code', 'some_value']
outputDF = spark.createDataFrame(data=output_li, schema=otherColumns)
outputDF.printSchema()
outputDF.show(truncate=False)

输出:

root
 |-- real_aid: string (nullable = true)
 |-- code: string (nullable = true)
 |-- some_value: string (nullable = true)

+--------+----+----------+
|real_aid|code|some_value|
+--------+----+----------+
|abc     |111 |desc111   |
|abc     |112 |desc112   |
|abc     |A123|city      |
|abc     |A456|state     |
|abc     |A789|zip       |
|xyz     |A123|city      |
|xyz     |A456|state     |
|xyz     |A789|zip       |
+--------+----+----------+

解决方案(基于Join实现)

可以通过构建关联映射表+Join的方式实现,适合大数据场景,具体步骤如下:

1. 生成关联映射表

从tempDF中提取需要的映射关系:每个aid需要关联自身和对应的bid的数据,每个bid只关联自身的数据。通过union生成映射表并去重:

# 生成aid到自身的映射
aid_self_map = tempDF.selectExpr("aid as target_aid", "aid as source_aid")
# 生成aid到bid的映射
aid_bid_map = tempDF.selectExpr("aid as target_aid", "bid as source_aid")
# 生成bid到自身的映射
bid_self_map = tempDF.selectExpr("bid as target_aid", "bid as source_aid")

# 合并映射表并去重,避免重复行
mapping_df = aid_self_map.union(aid_bid_map).union(bid_self_map).distinct()
mapping_df.show(truncate=False)

映射表输出:

+----------+-----------+
|target_aid|source_aid |
+----------+-----------+
|abc       |abc        |
|abc       |xyz        |
|xyz       |xyz        |
+----------+-----------+

2. 关联映射表与otherDF

将映射表与otherDF做内连接,通过source_aid和real_aid关联,最终将target_aid作为结果的real_aid:

from pyspark.sql.functions import broadcast

# 大数据场景下,广播小的映射表,避免大表shuffle,提升效率
result_df = broadcast(mapping_df).join(
    otherDF,
    mapping_df.source_aid == otherDF.real_aid,
    "inner"
).selectExpr(
    "target_aid as real_aid",
    "code",
    "some_value"
)

result_df.show(truncate=False)

运行后得到的结果与预期完全一致。

高效性说明

  • 广播小映射表(broadcast(mapping_df))可避免大表otherDF的shuffle操作,大幅提升Join效率。
  • 映射表通过distinct()去重,避免生成重复关联结果,减少不必要的数据处理。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 08:15:36