如何动态选择Spark DataFrame列:空值时自动取对应_p列值
动态处理Spark DataFrame的空值替换(年份列与对应_p列)
核心思路
先自动识别所有目标年份列(不带_p后缀的年份命名列),再对每个年份列使用coalesce函数实现"原列非空则取原列,原列为空则取对应_p列"的逻辑,全程无需硬编码列名。
Python 实现代码
import re from pyspark.sql.functions import coalesce, col # 从DataFrame列名中提取所有4位数字格式的年份列(如2019、2020) year_cols = [col_name for col_name in df.columns if re.match(r'^\d{4}$', col_name)] # 为每个年份列构造空值替换表达式,处理后保留原年份列名 processed_columns = [coalesce(col(c), col(f"{c}_p")).alias(c) for c in year_cols] # 生成最终DataFrame:若需保留其他非年份/非_p列,可添加到select列表中 # 示例:other_cols = [c for c in df.columns if not (c in year_cols or c.endswith('_p'))] # final_df = df.select(other_cols + processed_columns) final_df = df.select(processed_columns)
Scala 实现代码
import org.apache.spark.sql.functions.{coalesce, col} import scala.util.matching.Regex // 定义年份列匹配规则:4位数字 val yearPattern: Regex = """^\d{4}$""".r // 筛选出所有符合规则的年份列 val yearCols = df.columns.filter(yearPattern.findFirstIn(_).isDefined) // 构造每个年份列的空值替换逻辑 val processedCols = yearCols.map(c => coalesce(col(c), col(s"${c}_p")).alias(c)) // 生成最终DataFrame,如需保留其他列可自行补充到select参数中 val finalDf = df.select(processedCols: _*)
灵活调整说明
如果你的年份列命名不是严格的4位数字(比如带前缀sales_2019),只需修改正则表达式即可适配,例如将匹配规则改为r'.*(\d{4})$'来提取末尾的年份部分,再对应拼接_p后缀。
内容的提问来源于stack exchange,提问作者sparc
相关产品推荐
相关产品推荐

