【问题标题】:Is it possible to do filebased processing with UDF in pyspark?是否可以在 pyspark 中使用 UDF 进行基于文件的处理?
【发布时间】:2020-11-21 05:47:51
【问题描述】:

我定义了一个 UDF,它使用一个数据框执行以下操作,其中一列包含 zip 文件在 azure blob 存储中的位置(我在没有 spark 的情况下测试了 UDF,结果成功):

  1. 从 blob 下载定义的文件并将其安全地保存在执行器/驱动程序的某个位置
  2. 提取 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


    【解决方案1】:

    PySpark 读取器/写入器可以轻松地并行读取和写入文件。在 Spark 中工作时,通常不应在驱动程序节点上循环文件或保存数据。

    假设您在 my-bucket/my-folder 目录中有 100 个压缩后的 CSV 文件。以下是如何将它们并行读入 DataFrame:

    df = spark.read.csv("my-bucket/my-folder")
    

    以下是如何将它们写入 50 个 Snappy 压缩 Parquet 文件(并行):

    df.repartition(50).write.parquet("my-bucket/another-folder")
    

    读者/作者为您完成所有繁重的工作。有关repartition 的更多信息,请参阅here。

    【讨论】:

    • 那么用 spark.read.csv("my-bucket/my-folder") 会直接解压吗?所以不幸的是它是一个相当复杂的结构。 csv 不仅直接在压缩文件中,而且在 zip 中还有一个压缩文件,其中有多个文件夹和不同的文件,只有一个可以用于分析,应该读入数据框。
    猜你喜欢
    • 2012-01-19
    • 2023-02-25
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-11-09
    • 1970-01-01
    相关资源
    最近更新 更多