PySpark执行a.take(1)报PicklingError错误的排查求助
问题:PySpark执行
take(1)时出现PicklingError错误 环境信息
- Spark版本:spark-3.3.1-bin-hadoop3.tgz
- Java版本:Java(TM) SE Runtime Environment (build 1.8.0_351-b10)
- Python版本:3.11.1
操作过程
在CMD中启动PySpark并执行操作:
C:\Users\Administrator>SUCCESS: The process with PID 5328 (child process of PID 4476) has been terminated. SUCCESS: The process with PID 4476 (child process of PID 1092) has been terminated. SUCCESS: The process with PID 1092 (child process of PID 3952) has been terminated. pyspark Python 3.11.1 (tags/v3.11.1:a7a450f, Dec 6 2022, 19:58:39) [MSC v.1934 64 bit (AMD64)] on win32 Type "help", "copyright", "credits" or "license" for more information. Setting default log level to "WARN". To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel). 23/01/08 20:07:53 WARN NativeCodeLoader: Unable to load native-hadoop library for your platform... using builtin-java classes where applicable Welcome to ____ __ / __/__ ___ _____/ /__ _\ \/ _ \/ _ `/ __/ '_/ /__ / .__/\_,_/_/ /_/\_\ version 3.3.1 /_/ Using Python version 3.11.1 (tags/v3.11.1:a7a450f, Dec 6 2022 19:58:39) Spark context Web UI available at http://Mohit:4040 Spark context available as 'sc' (master = local[*], app id = local-1673188677388). SparkSession available as 'spark'. >>> 23/01/08 20:08:10 WARN ProcfsMetricsGetter: Exception when trying to compute pagesize, as a result reporting of ProcessTree metrics is stopped a = sc.parallelize([1,2,3,4,5,6,7,8,9,10])
执行a.take(1)时抛出错误,预期结果应为[1]。
错误信息
>>> a.take(1) Traceback (most recent call last): File "C:\Spark\python\pyspark\serializers.py", line 458, in dumps return cloudpickle.dumps(obj, pickle_protocol) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 73, in dumps cp.dump(obj) File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 602, in dump return Pickler.dump(self, obj) ^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 692, in reducer_override return self._function_reduce(obj) ^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 565, in _function_reduce return self._dynamic_function_reduce(obj) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 546, in _dynamic_function_reduce state = _function_getstate(func) ^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 157, in _function_getstate f_globals_ref = _extract_code_globals(func.__code__) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle.py", line 334, in _extract_code_globals out_names = {names[oparg]: None for _, oparg in _walk_global_ops(co)} ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle.py", line 334, in <dictcomp> out_names = {names[oparg]: None for _, oparg in _walk_global_ops(co)} ~~~~~^^^^^^^ IndexError: tuple index out of range Traceback (most recent call last): File "C:\Spark\python\pyspark\serializers.py", line 458, in dumps return cloudpickle.dumps(obj, pickle_protocol) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 73, in dumps cp.dump(obj) File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 602, in dump return Pickler.dump(self, obj) ^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 692, in reducer_override return self._function_reduce(obj) ^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 565, in _function_reduce return self._dynamic_function_reduce(obj) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 546, in _dynamic_function_reduce state = _function_getstate(func) ^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle_fast.py", line 157, in _function_getstate f_globals_ref = _extract_code_globals(func.__code__) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle.py", line 334, in _extract_code_globals out_names = {names[oparg]: None for _, oparg in _walk_global_ops(co)} ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\cloudpickle\cloudpickle.py", line 334, in <dictcomp> out_names = {names[oparg]: None for _, oparg in _walk_global_ops(co)} ~~~~~^^^^^^^ IndexError: tuple index out of range During handling of the above exception, another exception occurred: Traceback (most recent call last): File "<stdin>", line 1, in <module> File "C:\Spark\python\pyspark\rdd.py", line 1883, in take res = self.context.runJob(self, takeUpToNumLeft, p) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\context.py", line 1486, in runJob sock_info = self._jvm.PythonRDD.runJob(self._jsc.sc(), mappedRDD._jrdd, partitions) ^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\rdd.py", line 3505, in _jrdd wrapped_func = _wrap_function( ^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\rdd.py", line 3362, in _wrap_function pickled_command, broadcast_vars, env, includes = _prepare_for_python_RDD(sc, command) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\rdd.py", line 3345, in _prepare_for_python_RDD pickled_command = ser.dumps(command) ^^^^^^^^^^^^^^^^^^ File "C:\Spark\python\pyspark\serializers.py", line 468, in dumps raise pickle.PicklingError(msg) _pickle.PicklingError: Could not serialize object: IndexError: tuple index out of range
排查方案
- 核心原因:Spark 3.3.1官方支持的Python版本为3.7-3.10,Python 3.11不在支持列表内,导致自带的cloudpickle无法正确处理Python 3.11的字节码,引发序列化时的索引越界错误。Google Colab中使用的Python版本符合Spark兼容要求,因此无异常。
- 解决办法:
- 降级Python到3.10版本,完全适配Spark 3.3.1,这是最稳定的方案。
- 更新Spark内置的cloudpickle:找到Spark安装目录下的
python/pyspark/cloudpickle文件夹,替换为支持Python 3.11的最新版cloudpickle(需注意版本兼容性)。 - 升级Spark到3.4及以上版本,官方从Spark 3.4开始正式支持Python 3.11。
内容的提问来源于stack exchange,提问作者Mohit Aswani
相关产品推荐
相关产品推荐

