如何在PySpark DataFrame中应用列内表达式处理另一列?
在PySpark DataFrame中动态应用列存储的表达式处理另一列
问题描述
现有如下PySpark DataFrame,需要用expr_to_apply列中存储的Spark SQL表达式,对new_feed_dt列进行处理:
new_feed_dt regex_to_apply expr_to_apply 053021 | _(\d+) | date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd') 053022 | _(\d+) | date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd') 053023 | _(\d+) | date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd') 053024 | [a-zA-Z]+(\d+) | date_format(to_date(new_feed_dt, 'MMddyyyy'), 'yyyy-MM-dd') 053025 | DT(\d+) | date_format(to_date(new_feed_dt, 'MMddyy'), 'yyyy-MM-dd')
尝试了以下代码但未生效:
df_with_regex.withColumn( 'new_feed_dt', f.expr("expr_to_apply") )
原代码问题分析
f.expr("expr_to_apply")的作用是直接引用expr_to_apply列的字符串值,而非解析并执行这些字符串作为Spark SQL表达式。因此执行后,new_feed_dt列会被替换成原expr_to_apply列的表达式文本,而不是处理后的日期结果。
解决方案
由于每行的表达式可能不同,需要动态解析执行每行的表达式,以下提供几种可行方法:
方法一:使用RDD Map逐行执行表达式
通过将DataFrame转换为RDD,逐行替换表达式中的new_feed_dt为当前行的实际值,再通过Spark SQL执行表达式:
from pyspark.sql import SparkSession, Row import pyspark.sql.functions as f # 初始化SparkSession spark = SparkSession.builder.appName("DynamicExprApply").getOrCreate() # 构建示例DataFrame data = [ ("053021", "_(\d+)", "date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd')"), ("053022", "_(\d+)", "date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd')"), ("053023", "_(\d+)", "date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd')"), ("053024", "[a-zA-Z]+(\d+)", "date_format(to_date(new_feed_dt, 'MMddyyyy'), 'yyyy-MM-dd')"), ("053025", "DT(\d+)", "date_format(to_date(new_feed_dt, 'MMddyy'), 'yyyy-MM-dd')") ] df = spark.createDataFrame(data, ["new_feed_dt", "regex_to_apply", "expr_to_apply"]) # 定义行处理函数 def process_row(row): # 替换表达式中的列名为当前行的实际值(加单引号转义字符串) expr_str = row.expr_to_apply.replace("new_feed_dt", f"'{row.new_feed_dt}'") # 执行表达式并获取结果 processed_date = spark.sql(f"SELECT {expr_str} AS result").first().result # 返回新的Row对象 return Row( new_feed_dt=row.new_feed_dt, regex_to_apply=row.regex_to_apply, expr_to_apply=row.expr_to_apply, processed_new_feed_dt=processed_date ) # 转换为RDD处理后再转回DataFrame result_df = df.rdd.map(process_row).toDF() result_df.show(truncate=False)
方法二:使用Pandas UDF批量处理(更高效)
通过Pandas UDF批量处理分区数据,减少Spark SQL调用次数,提升效率:
from pyspark.sql.functions import pandas_udf import pandas as pd # 定义Pandas UDF处理函数 def apply_expr_batch(df: pd.DataFrame) -> pd.DataFrame: def apply_single_expr(row): expr_str = row["expr_to_apply"].replace("new_feed_dt", f"'{row['new_feed_dt']}'") return spark.sql(f"SELECT {expr_str} AS res").first().res # 对每行应用表达式 df["processed_new_feed_dt"] = df.apply(apply_single_expr, axis=1) return df[["new_feed_dt", "regex_to_apply", "expr_to_apply", "processed_new_feed_dt"]] # 注册Pandas UDF(指定输出Schema) apply_expr_udf = pandas_udf( apply_expr_batch, schema="new_feed_dt string, regex_to_apply string, expr_to_apply string, processed_new_feed_dt string" ) # 应用UDF处理数据 result_df = df.groupBy().apply(apply_expr_udf) result_df.show(truncate=False)
方法三:分支判断(适用于表达式种类有限的场景)
如果expr_to_apply的取值类型较少,直接用when...otherwise分支判断调用对应表达式,性能最优:
result_df = df.withColumn( "processed_new_feed_dt", f.when(f.col("expr_to_apply") == "date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd')", f.expr("date_format(to_date(new_feed_dt, 'yyyyMMdd'), 'yyyy-MM-dd')")) .when(f.col("expr_to_apply") == "date_format(to_date(new_feed_dt, 'MMddyyyy'), 'yyyy-MM-dd')", f.expr("date_format(to_date(new_feed_dt, 'MMddyyyy'), 'yyyy-MM-dd')")) .when(f.col("expr_to_apply") == "date_format(to_date(new_feed_dt, 'MMddyy'), 'yyyy-MM-dd')", f.expr("date_format(to_date(new_feed_dt, 'MMddyy'), 'yyyy-MM-dd')")) )
注意事项
- 前两种方法仅适用于表达式完全可控的场景,避免执行不可信的表达式(防止代码注入风险)。
内容的提问来源于stack exchange,提问作者Tomás Jullier
相关产品推荐
相关产品推荐

