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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 03:29:52