【发布时间】:2020-11-21 05:47:51
【问题描述】:
我定义了一个 UDF,它使用一个数据框执行以下操作,其中一列包含 zip 文件在 azure blob 存储中的位置(我在没有 spark 的情况下测试了 UDF,结果成功):
- 从 blob 下载定义的文件并将其安全地保存在执行器/驱动程序的某个位置
- 提取 blob 的某个文件并将其安全地保存在 Excutor/Driver 上
使用这个 UDF,我体验到它的速度就像我只是在 python 中循环文件一样。那么在 Spark 中甚至有可能完成这种任务吗?我想使用 spark 来并行化下载和解压缩以加快速度。 我通过 ssh 连接到 Excutor 和 Driver(它是一个测试集群,所以每个集群只有一个),发现只有数据在 Excutor 上处理,而驱动程序根本没有做任何事情。为什么会这样?
下一步是将提取的文件(普通 csvs)读取到 spark 数据帧。但是,如果文件分布在执行器和驱动程序上,如何做到这一点?我还没有找到访问执行器存储的方法。或者是否有可能在 UDF 中定义一个公共位置以将其写回到驱动程序的某个位置?
我想读取比提取的文件:
data_frame = (
spark
.read
.format('csv')
.option('header', True)
.option('delimiter', ',')
.load(f"/mydriverpath/*.csv"))
如果有其他方法可以并行下载和解压缩文件,我很乐意听到。
【问题讨论】:
标签: python apache-spark pyspark azure-blob-storage