【问题标题】:Spark: How do I write individual row to S3 / HDFS as JPGSpark:如何将单个行以 JPG 格式写入 S3/HDFS
【发布时间】:2021-05-22 22:56:50
【问题描述】:

我必须将数据作为单个 JPG 文件(约数百万)从 PySpark 写入 S3 存储桶。

我尝试了多种选择:

设置:AWS EMR 集群和 Jupyter 笔记本。

  1. 在 'foreach' 方法中创建一个 boto3 客户端并写入 S3 ==> 太慢且效率低下,因为我们为每个任务打开客户端。

def get_image(y):
    res = requests.get(img_url, stream=True)
    file_name = "./" +str(cid) + ".jpg"
    client = boto3.client('s3')
    file_name = str(cid) + ".jpg"
    client.put_object(Body=res.content, Bucket='test',  Key='out_images/'+file_name)

myRdd.foreach(get_image)
  1. 写入本地文件系统并运行“aws S3 复制”到 S3 => 如果将这些数据写入每个单独的工作节点的卷,则不清楚如何访问这些数据。在作业运行时登录到工作节点,但无法准确找到 JPG 的写入位置。

def get_image(y):
    res = requests.get(img_url, stream=True)
    file_name = "./" +str(cid) + ".jpg"
    with open(file_name, 'wb') as f:
        f.write(res.content)

myRdd.foreach(get_image)
  1. 写入 HDFS 并稍后运行 s3-dist-cp。可能是最有效的,但尚未在代码方面取得成功。 I get path cannot be found exceptions

def get_image(y):
    res = requests.get(img_url, stream=True)
    file_name = "hdfs://" +str(cid) + ".jpg"
    with open(file_name, 'wb') as f:
        f.write(res.content)

myRdd.foreach(get_image)

有人可以提出一个实现这一目标的好方法吗?

【问题讨论】:

  • 如果按行数分区然后以这种方式写入 S3 会怎样?我想您可以使用使用 boto3 的脚本来更改每个文件的格式。 rows = df.count(), df.repartition(rows).write.avro('save-dir')
  • 我有 5 亿张图片要写。我不认为重新分区是一个理想的解决方案。

标签: apache-spark pyspark amazon-emr


【解决方案1】:

如果将 foreach 替换为 foreachPartition,则解决方案 1 效果很好。此更改后,每个分区仅创建一个客户端:

def get_image(y_it):
    client = boto3.client('s3')
    for y in y_it:
        img_url = ...
        cid = ...
        res = requests.get(img_url, stream=True)
        file_name = str(cid) + ".jpg"
        client.put_object(Body=res.content, Bucket='test',  Key='out_images/'+file_name)

myRdd.foreachPartition(get_image)

y_it 的循环内,重复使用了同一个客户端。

如果将requests.Sessions 用于this answer 中所述的http 调用,事情甚至会变得更快。在这种情况下,会在循环外通过y_it(如客户端)创建一个 http 会话,然后在循环内重复使用。

【讨论】:

  • 感谢您的回复。我尝试了一个样本集,与方法 1 相比它肯定更快。我将在完整数据集上再试一次。只是好奇这是否是唯一的方法,我想听听您对方法 3 的看法。
  • 我会坚持方法 1,因为这种方法可以处理 Spark 中的所有内容,并且您不需要任何外部工具。如果您想在get_image 内将数据保存到 HDFS,您可以查看this answer。它是用 Java 编写的,但这个想法也应该适用于 Python
猜你喜欢
  • 2017-09-24
  • 2016-06-29
  • 2019-09-23
  • 2021-04-19
  • 2020-04-24
  • 1970-01-01
  • 2020-09-15
  • 2018-06-02
  • 2021-04-20
相关资源
最近更新 更多