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

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)
...

试过相关方法,结果仍报错。


解决方案

问题根源

  1. 格式字符串与实际行长度不匹配:你的FMT_STRING='3sx3sx4s3s7sx'计算总长度为23,但示例中account字段实际是13位,导致struct.unpack抛出异常,直接引发Python Worker崩溃。
  2. Row字段名与Schema不匹配:UDF返回的Row第一个字段名是tye,但Schema中定义的是type,字段名不匹配会导致数据解析错误。

修复步骤

  1. 修正格式字符串:根据实际行结构调整格式字符串,示例中account为13位,修正后格式字符串应为3sx3sx4s3s13sx。
  2. 匹配Row字段名与Schema:将返回Row的字段名从tye改为type,与Schema定义保持一致。
  3. 添加异常处理:在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 15:15:05