PySpark如何将DataFrame子集的列值传播填充到另一DataFrame
PySpark 匹配id填充固定列值实现方案
原写法的问题
你之前的代码没达到预期,核心是两个问题:
when判断没加otherwise分支,虽然你给的示例里A表不匹配id的Place本来就是null,但只要A表原有Place存在非null值,就会被错误覆盖成null;而且你直接硬写了'Office',如果B表的固定Place值换成别的(比如仓库Warehouse),还得手动改代码,适配性很差。- 用
collect把B表的id全拉到Driver端再走isin匹配,B表数据量稍微大一点就会把Driver内存撑爆,大列表的isin匹配性能也比Spark原生join差很多。
推荐方案(无硬编码、大数据量可用)
因为B表的Place列全是同一个固定值,直接走左关联就行,不用提前收集id,也不用硬编码Place的取值,全程分布式执行性能最好:
from pyspark.sql import functions as F # 对B表按id去重,避免id重复导致关联后数据行数膨胀 B_slim = B.select("id", "Place").dropDuplicates(["id"]) # 左关联,匹配到的行取B表的Place值,匹配不到保留A表原有的Place值 C = A.alias("a").join( B_slim.alias("b"), on="id", how="left" ).select( "id", "Name", F.coalesce(F.col("b.Place"), F.col("a.Place")).alias("Place") )
这个写法不管B表的固定Place是Office还是Warehouse,都不用改代码,自动适配。
小数据量快速修复版
如果B表数据量很小(几千条以内),你想沿用之前的isin逻辑,只要补上otherwise分支,再动态取B表的固定Place值就行,不用硬编码:
from pyspark.sql import functions as F # 自动获取B表的固定Place值,不用手动填 fixed_place_val = B.select("Place").first()[0] B_id_list = [row.id for row in B.select("id").collect()] C = A.withColumn( "Place", F.when(F.col("id").isin(B_id_list), fixed_place_val).otherwise(F.col("Place")) )
注意:这个方案只适合B表数据量极小的场景,B表数据超过10万条不要用collect操作,会有性能和内存问题。
内容的提问来源于stack exchange,提问作者tom.bell
相关产品推荐
相关产品推荐

