Spark collect()在Jupyter中触发IllegalArgumentException问题求助
我之前碰到过类似的场景,结合你的环境信息,给你几个不用降级Spark也能尝试排查的方向:
Anaconda环境的依赖版本冲突
PySpark 2.2.1依赖特定版本的py4j(记得是0.10.4左右),而Anaconda默认自带的py4j版本可能和Spark内置的不一致。Jupyter运行在Anaconda的Python环境中,会优先使用Anaconda的py4j,但pyspark客户端启动时会自动加载Spark自带的py4j,这就导致了驱动端和executor端的通信序列化逻辑不匹配——collect()需要把所有数据拉回驱动,对序列化的一致性要求更高,而take()只拉取少量数据,刚好没触发问题。
你可以分别在Jupyter和pyspark客户端里运行import py4j; print(py4j.__version__)对比版本,如果不一致,要么把Anaconda的py4j降级/升级到Spark对应的版本,要么在Jupyter的配置里指定优先加载Spark自带的py4j路径。Jupyter内核的环境变量配置缺失
pyspark客户端启动时会自动设置一系列Spark相关的环境变量(比如SPARK_HOME、PYSPARK_PYTHON、PYSPARK_DRIVER_MEMORY等),但Jupyter的内核可能没正确继承这些变量。比如如果驱动端内存配置不足,collect()拉取大量数据时会触发内存溢出,而take()因为数据量小不会;或者PYSPARK_PYTHON配置不一致,导致executor用的Python环境和驱动端不同,序列化数据时出错。
你可以在Jupyter里运行import os; print(os.environ.get('SPARK_HOME'))和print(os.environ.get('PYSPARK_PYTHON')),对比pyspark客户端里的输出,确保变量一致。如果缺失,可以在Jupyter的内核配置文件里手动添加这些环境变量。Jupyter命名空间的序列化问题
如果你在Jupyter里定义了自定义的函数、类或者修改了全局变量,这些内容的序列化逻辑和pyspark客户端里的可能不一样。collect()需要将executor端处理后的数据反序列化为驱动端的对象,如果驱动端的类定义和executor端不匹配(比如Jupyter里的类被重新定义过),就会触发反序列化错误,而take()可能刚好没碰到包含这些自定义类型的数据。
可以尝试在Jupyter里重新启动内核,只运行必要的Spark代码(不定义额外的自定义对象),再测试collect()是否正常。Spark 2.2.1的特定bug(可尝试workaround)
虽然你不想降级,但有些Spark 2.2.1的bug有临时解决办法。比如当DataFrame包含嵌套复杂类型(如数组、结构体)时,collect()在Jupyter中可能触发序列化异常,这时可以先将DataFrame转换为RDD再调用collect(),或者用toPandas()替代(如果数据量不大的话),这些方法可能绕过原有的bug。
内容的提问来源于stack exchange,提问作者Rodrygo Santos

