PySpark中UDF解码cp1047失败,请求故障原因排查建议
Python环境编码支持不一致:本地Python环境能识别cp1047,要么是系统预装了对应字符集包,要么依赖了第三方库,但PySpark Worker节点的Python环境缺少这些支持。cp037是Python标准库默认自带的EBCDIC编码,而cp1047不在默认支持列表中,Worker环境缺依赖就会触发编码找不到的错误。
集群系统字符集缺失:如果是Linux集群,默认的glibc字符集包可能未包含cp1047的定义。本地机器可能已经安装了全量字符集,但Worker节点仅安装了基础字符集,导致Python无法识别cp1047编码。
Python版本不匹配:老版本Python(比如3.6及更早)对EBCDIC编码的支持有限,未内置cp1047的编码规则。你本地使用的高版本Python能正常解码,但Spark集群使用的老版本Python不支持该编码。
核对Spark Worker节点与本地的Python版本,若集群版本过低,升级到Python 3.7及以上版本——新版本对更多EBCDIC编码有内置支持。
在所有Worker节点安装全量字符集包:RedHat/CentOS系统执行
yum install glibc-langpack-*,Debian/Ubuntu系统执行apt-get install locales-all,安装完成后重启Spark服务。若无法修改集群环境,可借助第三方编码库处理:先在所有Worker节点安装
ebcdic库(执行pip install ebcdic),再修改UDF代码如下:
from ebcdic import cp1047 from pyspark.sql.functions import udf from pyspark.sql.types import StringType def decode_cp1047(input_bytes): return cp1047.decode(input_bytes) decode_cp1047_udf = udf(decode_cp1047, StringType())
内容的提问来源于stack exchange,提问作者Kaushik Ghosh

