基于两个DataFrame,使用when函数生成含全量SN列表列的技术问询
Spark DataFrame条件生成列表列的实现方案
需求说明
现有两个DataFrame:
- df1:包含SN列(序列号)
- df2:包含CONCERN列
需要给df2新增一列EFFECTIVITY,规则为:
- 当df2的CONCERN列包含
all(不区分大小写)时,填入df1中SN列所有序列号组成的列表 - 否则填空字符串
原代码问题分析
原代码存在以下几个问题:
df1.select(collect_list("SN")).show()执行后all_data的值为None,因为show()仅用于打印DataFrame内容,不会返回实际数据- 代码中
df = df.withColumn(...)的df未定义,实际应该操作目标DataFramedf2 - 直接使用
df2.CONCERN.contains('ALL')会导致大小写敏感的匹配问题,且未正确引用列对象
正确实现代码
步骤1:导入依赖函数
from pyspark.sql import functions as F
步骤2:提取df1的SN全量列表
# 从df1中提取所有SN值组成Python列表 all_sn_list = df1.select(F.collect_list("SN")).first()[0]
步骤3:给df2新增EFFECTIVITY列
# 处理大小写匹配,将Python列表转为Spark ArrayType列 df2_result = df2.withColumn( "EFFECTIVITY", F.when( F.lower(F.col("CONCERN")).contains("all"), F.array(*[F.lit(sn) for sn in all_sn_list]) ).otherwise(F.lit("")) )
性能优化方案(大数据量场景)
如果df1的SN数据量较大,使用广播变量减少节点间数据传输:
from pyspark.sql import SparkSession # 获取SparkSession实例 spark = SparkSession.builder.getOrCreate() # 广播全量SN列表 broadcast_all_sn = spark.sparkContext.broadcast(all_sn_list) # 生成结果DataFrame df2_result = df2.withColumn( "EFFECTIVITY", F.when( F.lower(F.col("CONCERN")).contains("all"), F.array(*[F.lit(sn) for sn in broadcast_all_sn.value]) ).otherwise(F.lit("")) )
内容的提问来源于stack exchange,提问作者SISI
相关产品推荐
相关产品推荐

