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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 18:13:27