【问题标题】:Reading multiple files from different aws S3 in Spark parallelly在 Spark 中并行读取来自不同 aws S3 的多个文件
【发布时间】:2023-01-24 11:01:17
【问题描述】:

我有一个场景,我需要从位于不同位置和不同架构的 s3 存储桶中读取许多文件(csv 或镶木地板)。

我这样做的目的是从不同的 s3 位置提取所有元数据信息并将其保存为 Dataframe 并将其另存为 s3 本身中的 csv 文件。这里的问题是我有很多 s3 位置来读取文件(分区)。我的示例 s3 位置就像

s3://myRawbucket/source1/filename1/year/month/day/16/f1.parquet
s3://myRawbucket/source2/filename2/year/month/day/16/f2.parquet
s3://myRawbucket/source3/filename3/year/month/day/16/f3.parquet
s3://myRawbucket/source100/filename100/year/month/day/16/f100.parquet
s3://myRawbucket/source150/filename150/year/month/day/16/f150.parquet    and .......... so on

我需要做的就是使用 spark 代码读取这么多文件(大约 200 个)并根据需要应用一些转换并提取标头信息、计数信息、s3 位置信息、数据类型。

读取所有这些文件(不同模式)并使用火花代码(Dataframe)处理它并将其保存为 s3 存储桶中的 csv 的有效方法是什么?请耐心等待,因为我是 Spark 世界的新手。我正在使用 python (Pyspark)

【问题讨论】:

  • 您可以尝试 multiprocessing / Thread 并行处理文件。
  • 据我所知,spark 用于并行处理。我如何使用 spark 实现它?

标签: python apache-spark pyspark apache-spark-sql boto3


【解决方案1】:

我想你想要做的是使用一些 Python/Pandas 逻辑并使用 Spark 并行化作业。 Fugue 非常适合。您可以通过极少的代码更改将您的逻辑移植到 Spark。让我们先担心用 Python 和 Pandas 定义逻辑,然后我们可以把它带到 Spark 中。

首先是设置:

import pandas as pd

df = pd.DataFrame({"x": [1,2,3]})
df.to_parquet("/tmp/1.parquet")
df.to_parquet("/tmp/2.parquet")
df.to_parquet("/tmp/3.parquet")

我们需要一个包含所有文件的小型 DataFrame 来使用 Spark 编排作业。例如:

file_paths = pd.DataFrame({"path": ["/tmp/1.parquet",
                                    "/tmp/2.parquet",
                                    "/tmp/3.parquet"]})

现在我们可以创建一个函数来保存每个文件的逻辑。请注意,当我们将其引入 Spark 时,我们将为每个文件路径创建 1 个“作业”。我们的函数一次只需要能够处理一个文件。

def process(df:pd.DataFrame) -> pd.DataFrame:
    path = df.iloc[0]['path']
    
    tmp = pd.read_parquet(path)
    
    # transformation
    tmp['y'] = tmp['x'] + 1
    
    # save
    tmp.to_parquet(path)
    
    # summary stats
    return pd.DataFrame({"path": [path],
                         'count': [tmp.shape[0]]})

我们可以测试代码:

process(file_paths)

这给了我们:

path    count
/tmp/1.parquet  3

现在我们可以使用 Fugue 将其引入 Spark。我们只需要 transform() 函数将逻辑引入 Spark。该模式是 Spark 的要求。

import fugue.api as fa
from pyspark.sql import SparkSession

spark = SparkSession.builder.getOrCreate()

out = fa.transform(file_paths, process, schema="path:str,count:int", engine=spark)

# out is a Spark DataFrame
out.show()

输出将是:

+--------------+-----+
|          path|count|
+--------------+-----+
|/tmp/1.parquet|    3|
|/tmp/2.parquet|    3|
|/tmp/3.parquet|    3|
+--------------+-----+

【讨论】:

    猜你喜欢
    • 2017-04-25
    • 2017-12-23
    • 1970-01-01
    • 2019-10-27
    • 1970-01-01
    • 2016-12-04
    • 1970-01-01
    • 1970-01-01
    • 2018-10-18
    相关资源
    最近更新 更多