PySpark多列返回UDF解析定宽文件报错:Python Worker无法连接
问题
尝试用PySpark读取定宽文件,计划通过Python的struct模块解析,定义调用struct.unpack的UDF处理后返回多列到原DataFrame,但运行报错。
文件内容示例:
CD EUR 0025800202878 CD EUR 0025800203059 CD EUR 0025800203070 CD EUR 0025800203085 CD EUR 0025800203110
使用的代码:
# 库导入和SparkSession初始化省略 FMT_STRING='3sx3sx4s3s7sx' schema = StructType([ StructField('type', StringType(), False), StructField('value', StringType(), False), StructField('group', LongType(), False), StructField('cat', LongType(), False), StructField('account', LongType(), False) ]) def unpacker(x): values = struct.unpack(FMT_STRING, x.encode('utf-8')) return Row('tye', 'value', 'group', 'cat', 'account')( values[0].decode('utf-8').strip(), values[1].decode('utf-8').strip(), int(values[2].decode('utf-8').strip()), int(values[3].decode('utf-8').strip()), int(values[4].decode('utf-8').strip()) ) unpacker_udf = F.udf(unpacker, schema) data = spark.read.format('text').load('example_file.txt') data = data.withColumn('results', unpacker_udf(F.col('value'))) data.show() # 另一种尝试: # data = spark.read.format('text').load('example_file.txt') # data2 = data.select(unpacker_udf('value').alias('results')) # data2.printSchema() # ... 能显示schema # data2.show() # ... 抛出相同异常
执行到show()时出现错误:
... Py4JJavaError: An error occurred while calling o128.showString. : org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in stage 6.0 failed 1 times, most recent failure: Lost task 0.0 in stage 6.0 (TID 6) (xyz.abc executor driver): org.apache.spark.SparkException: Python worker failed to connect back. at org.apache.spark.api.python.PythonWorkerFactory.createSimpleWorker(PythonWorkerFactory.scala:192) ...
试过相关方法,结果仍报错。
解决方案
问题根源
- 格式字符串与实际行长度不匹配:你的
FMT_STRING='3sx3sx4s3s7sx'计算总长度为23,但示例中account字段实际是13位,导致struct.unpack抛出异常,直接引发Python Worker崩溃。 - Row字段名与Schema不匹配:UDF返回的Row第一个字段名是
tye,但Schema中定义的是type,字段名不匹配会导致数据解析错误。
修复步骤
- 修正格式字符串:根据实际行结构调整格式字符串,示例中
account为13位,修正后格式字符串应为3sx3sx4s3s13sx。 - 匹配Row字段名与Schema:将返回Row的字段名从
tye改为type,与Schema定义保持一致。 - 添加异常处理:在UDF中加入try-except捕获解析错误,避免Worker直接崩溃,方便调试问题。
修复后的代码示例:
import struct from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType from pyspark.sql import functions as F from pyspark.sql import Row spark = SparkSession.builder.appName("FixedWidthReader").getOrCreate() # 修正后的格式字符串,匹配实际行长度 FMT_STRING='3sx3sx4s3s13sx' schema = StructType([ StructField('type', StringType(), False), StructField('value', StringType(), False), StructField('group', LongType(), False), StructField('cat', LongType(), False), StructField('account', LongType(), False) ]) def unpacker(x): try: values = struct.unpack(FMT_STRING, x.encode('utf-8')) return Row('type', 'value', 'group', 'cat', 'account')( values[0].decode('utf-8').strip(), values[1].decode('utf-8').strip(), int(values[2].decode('utf-8').strip()), int(values[3].decode('utf-8').strip()), int(values[4].decode('utf-8').strip()) ) except Exception as e: # 返回空值标记异常行,便于排查 return Row('type', 'value', 'group', 'cat', 'account')(None, None, None, None, None) unpacker_udf = F.udf(unpacker, schema) data = spark.read.format('text').load('example_file.txt') data = data.withColumn('results', unpacker_udf(F.col('value'))) # 展开嵌套列 data = data.select('value', 'results.*') data.show()
替代方案(无需UDF,更高效)
PySpark 3.0+原生支持定宽文件读取,无需自定义UDF:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType spark = SparkSession.builder.appName("FixedWidthReader").getOrCreate() schema = StructType([ StructField('type', StringType(), False), StructField('value', StringType(), False), StructField('group', LongType(), False), StructField('cat', LongType(), False), StructField('account', LongType(), False) ]) # 直接用csv格式读取定宽文件 data = spark.read.format("csv") \ .option("sep", None) \ .option("fixedWidth", "type:3,value:3,group:4,cat:3,account:13") \ .schema(schema) \ .load("example_file.txt") data.show()
内容的提问来源于stack exchange,提问作者nicodds
相关产品推荐
相关产品推荐

