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

在PySpark中运行自定义UDF时遭遇Py4JJavaError错误求助

Spark UDF执行错误排查与解决方案

问题场景

执行以下自定义UDF代码时触发Py4JJavaError:

from pyspark.sql.functions import udf
def retrieve_night_flag(row):
    x=row
    if x >= night_time['start_time'] and x < night_time['end_time']:
        return 'NIGHT'
    elif x >= day_time['start_time'] and x < day_time['end_time']:
        return 'DAY' 
    else:
        return 'EVE'

UDF_NAME = udf(lambda row: retrieve_night_flag(row),StringType())
df_source3 =df_source2 .withColumn('night_flag', UDF_NAME(col('hour_stamp')))

错误信息片段:

Py4JJavaError: An error occurred while calling o765.showString.
: org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 38.0 failed 1 times, most recent failure: Lost task 0.0 in stage 38.0 (TID 35) (INNOLX4356.itte.com executor driver): org.apache.spark.SparkException: 
Error from python worker:

      Traceback (most recent call last):
        File "/usr/lib/python3.5/runpy.py", line 174, in _run_module_as_main
          mod_name, mod_spec, code = _get_module_details(mod_name, _Error)
        File "/usr/lib/python3.5/runpy.py", line 109, in _get_module_details
          __import__(pkg_name)
        File "&lt;frozen importlib._bootstrap&gt;", line 969, in _find_and_load
        File "&lt;frozen importlib._bootstrap&gt;", line 958, in _find_and_load_unlocked
        File "&lt;frozen importlib._bootstrap&gt;", line 664, in _load_unlocked

错误原因

  • 外部变量序列化失败:night_time和day_time是Driver端定义的变量,Spark将UDF分发到Executor时需要序列化这些变量,若变量无法被正确序列化,会导致Python Worker加载失败。
  • UDF定义冗余:用lambda包裹retrieve_night_flag属于冗余写法,可能额外引入序列化逻辑问题。
  • Python版本兼容问题:错误日志显示使用Python3.5,部分Spark版本对Python3.5的兼容性较差,易引发Worker启动异常。

解决方案

方案1:修复UDF的变量传递方式

使用functools.partial将外部变量绑定到UDF函数,确保变量能被序列化传递:

from pyspark.sql.functions import udf, col
from pyspark.sql.types import StringType
from functools import partial

def retrieve_night_flag(night_time, day_time, x):
    if x >= night_time['start_time'] and x < night_time['end_time']:
        return 'NIGHT'
    elif x >= day_time['start_time'] and x < day_time['end_time']:
        return 'DAY' 
    else:
        return 'EVE'

# 绑定外部变量到函数
bound_udf = partial(retrieve_night_flag, night_time, day_time)
# 定义UDF
night_flag_udf = udf(bound_udf, StringType())
# 应用UDF到DataFrame
df_source3 = df_source2.withColumn('night_flag', night_flag_udf(col('hour_stamp')))

方案2:改用Spark内置函数替代UDF(推荐)

Spark内置函数性能优于自定义UDF,且无需处理变量序列化问题,直接用when实现逻辑:

from pyspark.sql.functions import col, when

df_source3 = df_source2.withColumn(
    'night_flag',
    when(
        (col('hour_stamp') >= night_time['start_time']) & (col('hour_stamp') < night_time['end_time']),
        'NIGHT'
    ).when(
        (col('hour_stamp') >= day_time['start_time']) & (col('hour_stamp') < day_time['end_time']),
        'DAY'
    ).otherwise('EVE')
)

方案3:升级Python版本

若使用的Spark版本较新,建议将Python升级至3.6及以上版本,避免版本不兼容导致的Worker启动失败。

内容的提问来源于stack exchange,提问作者SWARITRA SAXENA

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 01:48:19