如何在PySpark作业的Driver或Executor端列出wheel包中的Python包?
解决Spark中获取--py-files指定的Wheel包内容问题
一、Driver端解析Wheel文件
listJars()仅针对Java/Scala JAR包,要获取Python Wheel的内容,首先可以通过SparkContext获取所有--py-files上传的文件,再解析Wheel本身(本质是Zip格式包)。
方法1:用内置zipfile模块解析
import zipfile # 获取所有通过--py-files上传的文件路径 py_files = spark.sparkContext.getFiles() # 筛选出wheel文件 wheel_path = next((f for f in py_files if f.endswith('.whl')), None) if wheel_path: with zipfile.ZipFile(wheel_path, 'r') as zip_ref: # 读取METADATA文件获取包详细信息 metadata_paths = [name for name in zip_ref.namelist() if '/METADATA' in name] if metadata_paths: with zip_ref.open(metadata_paths[0]) as f: metadata_content = f.read().decode('utf-8') print("Wheel包元信息:") print(metadata_content) # 也可以直接列出包内所有文件 print("\nWheel包内所有文件:") for file_name in zip_ref.namelist(): print(file_name) else: print("未检测到通过--py-files上传的wheel文件")
方法2:用wheel库专业解析(需提前安装)
如果安装了wheel库,可以更便捷地提取包名、版本等结构化信息:
from wheel.wheelfile import WheelFile wheel_path = next((f for f in spark.sparkContext.getFiles() if f.endswith('.whl')), None) if wheel_path: with WheelFile(wheel_path) as wf: print(f"包名称:{wf.parsed_filename.group('name')}") print(f"包版本:{wf.parsed_filename.group('ver')}") print("\n包内文件列表:") for entry in wf.iterfiles(): print(entry.name)
二、Executor端解析Wheel文件
要在Executor节点上解析,需要把解析逻辑放到分布式任务中(比如mapPartitions),Spark会自动将--py-files的wheel文件分发到Executor的工作目录:
def parse_wheel_in_executor(_): import zipfile py_files = spark.sparkContext.getFiles() wheel_path = next((f for f in py_files if f.endswith('.whl')), None) result = [] if wheel_path: with zipfile.ZipFile(wheel_path, 'r') as zip_ref: metadata_paths = [name for name in zip_ref.namelist() if '/METADATA' in name] if metadata_paths: with zip_ref.open(metadata_paths[0]) as f: result.append(f.read().decode('utf-8')) else: result.append("Wheel包中未找到METADATA文件") else: result.append("Executor未接收到wheel文件") return result # 提交分布式任务执行解析 executor_results = spark.sparkContext.parallelize([1]).mapPartitions(parse_wheel_in_executor).collect() for res in executor_results: print("Executor端解析结果:\n", res)
三、为什么listJars()无法获取Wheel文件?
spark.sparkContext._jsc.sc().listJars()是JVM层面的API,仅用于枚举Spark作业加载的Java/Scala JAR包,而--py-files上传的Python文件(包括wheel)属于Python runtime的依赖,不会被JVM的Jar管理机制收录,因此无法通过该方法获取。
补充:本地提前查看Wheel内容
如果不需要在运行时解析,也可以在提交作业前,用命令行工具查看wheel内容:
# 查看包元信息 wheel show your_package.whl # 列出包内所有文件 unzip -l your_package.whl
内容的提问来源于stack exchange,提问作者Santosh Kumar Sahoo
相关产品推荐
相关产品推荐

