PySpark DataFrame应用UDF报错排查及代码去重方案咨询
问题描述
环境与数据
导入PySpark库的代码:
from pyspark.sql import SparkSession from pyspark import SparkConf from pyspark.sql import functions as F from pyspark.sql import types as T conf = SparkConf().set("spark.sql.catalogImplementation","hive") spark = SparkSession.builder.appName("SparkApp").config(conf=conf).getOrCreate() sc = spark.sparkContext
现有DataFrame df_delay 的数据样例:
+-----------+----------------------+----------------------+ |Call Number|Received DtTm |Dispatch DtTm | +-----------+----------------------+----------------------+ |1030101 |04/12/2000 09:00:29 PM|04/12/2000 09:02:00 PM| |1030104 |04/12/2000 09:09:02 PM|04/12/2000 09:10:29 PM| |1030106 |04/12/2000 09:09:44 PM|04/12/2000 09:11:47 PM| |1030107 |04/12/2000 09:13:47 PM|04/12/2000 09:14:13 PM| |1030108 |04/12/2000 09:14:43 PM|04/12/2000 09:16:24 PM| +-----------+----------------------+----------------------+
Schema信息:
root |-- Call Number: integer (nullable = true) |-- Received DtTm: string (nullable = true) |-- Dispatch DtTm: string (nullable = true)
错误尝试
尝试编写UDF将日期列转为Timestamp类型:
convert_to_datetime = F.udf(lambda l: F.unix_timestamp(l,"MM/dd/yyyy hh:mm:ss a").cast(T.TimestampType()), T.TimestampType())
应用UDF时触发错误:
df_delay.withColumn("DispatchTS", convert_to_datetime('Dispatch DtTm') ).show(5)
错误信息:
PythonException: An exception was thrown from the Python worker. Please see the stack trace below. Traceback (most recent call last): File "C:\Users\ajink\AppData\Local\Temp\ipykernel_9092\3273235857.py", line 1, in <lambda> File "C:\software\programming\spark-3.4.0-bin-hadoop3\python\lib\pyspark.zip\pyspark\sql\utils.py", line 159, in wrapped return f(*args, **kwargs) File "C:\software\programming\spark-3.4.0-bin-hadoop3\python\lib\pyspark.zip\pyspark\sql\functions.py", line 5254, in unix_timestamp return _invoke_function("unix_timestamp", _to_java_column(timestamp), format) File "C:\software\programming\spark-3.4.0-bin-hadoop3\python\lib\pyspark.zip\pyspark\sql\column.py", line 63, in _to_java_column jcol = _create_column_from_name(col) File "C:\software\programming\spark-3.4.0-bin-hadoop3\python\lib\pyspark.zip\pyspark\sql\column.py", line 55, in _create_column_from_name sc = get_active_spark_context() File "C:\software\programming\spark-3.4.0-bin-hadoop3\python\lib\pyspark.zip\pyspark\sql\utils.py", line 201, in get_active_spark_context raise RuntimeError("SparkContext or SparkSession should be created first.") RuntimeError: SparkContext or SparkSession should be created first.
可行但重复的实现
直接在DataFrame上应用逻辑实现了需求,但存在重复代码:
df_delay.withColumn("DispatchTS",F.unix_timestamp(F.col('Dispatch DtTm'), 'MM/dd/yyyy hh:mm:ss a').cast(T.TimestampType())) \ .withColumn("ReceivedTS",F.unix_timestamp(F.col('Received DtTm'), 'MM/dd/yyyy hh:mm:ss a').cast(T.TimestampType())) \ .withColumn('Difference',(F.col('DispatchTS')-F.col('ReceivedTS')).cast(T.IntegerType())/60) \ .filter(F.col('Difference') > 5) \ .show()
需求:找到避免重复代码的解决方案。
问题分析与解决方案
1. UDF报错原因
PySpark内置函数(如F.unix_timestamp)不能在Python UDF的lambda表达式中直接调用。UDF运行在独立的Python worker进程中,而Spark内置函数基于JVM执行,在UDF内部调用会导致无法获取当前SparkContext,从而触发报错。
正确思路:要么用纯Python日期处理逻辑写UDF(如datetime模块),要么直接使用Spark内置函数链(性能更优,避免跨进程序列化开销)。
2. 避免重复代码的几种方案
方案一:封装转换逻辑为自定义函数
封装生成转换逻辑的函数,传入列名即可复用:
def to_timestamp_col(col_name, fmt="MM/dd/yyyy hh:mm:ss a"): return F.unix_timestamp(F.col(col_name), fmt).cast(T.TimestampType()) # 应用到DataFrame df_result = df_delay.withColumn("DispatchTS", to_timestamp_col("Dispatch DtTm")) \ .withColumn("ReceivedTS", to_timestamp_col("Received DtTm")) \ .withColumn('Difference', (F.col('DispatchTS') - F.col('ReceivedTS')).cast(T.IntegerType())/60) \ .filter(F.col('Difference') > 5) df_result.show()
方案二:使用selectExpr批量处理
通过字符串表达式批量生成转换逻辑:
fmt = "MM/dd/yyyy hh:mm:ss a" select_expr = [ "*", f"cast(unix_timestamp(`Dispatch DtTm`, '{fmt}') as timestamp) as DispatchTS", f"cast(unix_timestamp(`Received DtTm`, '{fmt}') as timestamp) as ReceivedTS" ] df_result = df_delay.selectExpr(*select_expr) \ .withColumn('Difference', (F.col('DispatchTS') - F.col('ReceivedTS')).cast(T.IntegerType())/60) \ .filter(F.col('Difference') > 5) df_result.show()
方案三:循环处理多列
如果有大量日期列需要转换,用循环批量添加列:
fmt = "MM/dd/yyyy hh:mm:ss a" date_cols = ["Received DtTm", "Dispatch DtTm"] df_temp = df_delay for col_name in date_cols: new_col_name = col_name.replace(" DtTm", "TS") df_temp = df_temp.withColumn(new_col_name, F.unix_timestamp(F.col(col_name), fmt).cast(T.TimestampType())) df_result = df_temp.withColumn('Difference', (F.col('DispatchTS') - F.col('ReceivedTS')).cast(T.IntegerType())/60) \ .filter(F.col('Difference') > 5) df_result.show()
方案四:使用to_timestamp简化转换(Spark 2.2+)
Spark 2.2+提供F.to_timestamp函数,直接将字符串转为Timestamp类型,写法更简洁:
def to_timestamp_col(col_name, fmt="MM/dd/yyyy hh:mm:ss a"): return F.to_timestamp(F.col(col_name), fmt) df_result = df_delay.withColumn("DispatchTS", to_timestamp_col("Dispatch DtTm")) \ .withColumn("ReceivedTS", to_timestamp_col("Received DtTm")) \ .withColumn('Difference', (F.col('DispatchTS').cast(T.LongType()) - F.col('ReceivedTS').cast(T.LongType()))/60) \ .filter(F.col('Difference') > 5) df_result.show()
注:Timestamp直接相减得到Interval类型,需转为Long类型(毫秒时间戳)后计算差值,除以60得到分钟数。
内容的提问来源于stack exchange,提问作者AJ22
相关产品推荐
相关产品推荐

