如何用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
相关产品推荐
相关产品推荐

