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

基于两个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未定义,实际应该操作目标DataFrame df2
  • 直接使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 15:10:16