【问题标题】:Merge million S3 files generated hourly合并每小时生成的数百万个 S3 文件
【发布时间】:2021-08-20 05:45:29
【问题描述】:

我每小时创建数百万个文件。每个文件有一行数据。这些文件需要合并到一个文件中。

我尝试过以下方式:-

  1. 使用 aws s3 cp 下载一小时的文件。
  2. 使用 bash 命令合并文件。 或
  3. 使用 python 脚本合并文件。

这项每小时作业正在 Airflow on Kubernetes (EKS) 中运行。这需要一个多小时才能完成,并且正在创建积压工作。另一个问题是它经常导致 EC2 节点由于高 CPU 和内存使用率而停止响应。运行此作业的最有效方法是什么?

python脚本供参考:-

from os import listdir
import sys
# from tqdm import tqdm

files = listdir('./temp/')
dest = sys.argv[1]

data = []

tot_len = len(files)
percent = tot_len//100

for i, file in enumerate(files):
    if(i % percent == 0):
        print(f'{i/percent}% complete.')
    with open('./temp/'+file, 'r') as f:
        d = f.read()
        data.append(d)

result = '\n'.join(data)

with open(dest, 'w') as f:
    f.write(result)

【问题讨论】:

  • 有一种方法可以检查这个 - stackoverflow.com/questions/17749058/… 还有另一种有效且简单的方法,检查这个 - stackoverflow.com/questions/4827453/…
  • 显示代码?不清楚“合并”对您意味着什么,如果它只是意味着以任意顺序将文件内容粘贴在一起,那么很难想象 如何 编写需要大量内存使用的代码。
  • @user84634 我开始使用 cat,它在正则表达式扩展太多文件时出错。所以我用 find 命令替换了它。即使在此之后,该命令也需要很多时间。我尝试了读取然后加入的简单 python 代码,但即使这样也需要超过 2 小时,这对于每小时任务来说并不好,有时甚至会导致 EKS 集群中的节点崩溃。
  • @TimPeters 是的,我的意思是按任意顺序将文件粘贴在一起。当我运行“aws s3 cp”时,文件夹的内容被下载到 EKS 集群中运行的 pod 中,当我检查内存使用情况时,它开始超过 1GB 进入作业。 (原始文件的总大小约为 1GB)。我还添加了python代码。
  • 还有另一种方法,您可以编写一个脚本,将文件分成 50-100 个子目录。然后一次在 30-40 个目录中异步运行 python/bash 代码。与以前的方法相比,这将更快地完成工作。

标签: python amazon-s3 cron airflow amazon-eks


【解决方案1】:

把它放在那里以防其他人需要它。

我尽我所能优化了合并代码,但瓶颈仍然是读取或下载 s3 文件,即使使用官方的 aws cli 也很慢。

我发现了一个库 s5cmd,它非常快,因为它充分利用了多处理和多线程,它解决了我的问题。

链接:- https://github.com/peak/s5cmd

【讨论】:

    【解决方案2】:

    我希望您应该非常认真地考虑遵循您获得的 AWS 特定答案中的想法。我会将此响应添加为注释,除非无法在注释中连贯地显示缩进代码。

    对于您的 Python 脚本,您正在构建一个巨大的字符串,其字符数等于所有输入文件中的字符总数。所以当然内存使用量至少会增长到那么大。

    在读取文件内容后立即写出文件内容会占用更少的内存(请注意,此代码未经测试 - 可能有错字,我不知道):

    with open(dest, 'w') as fout:
        for i, file in enumerate(files):
            if(i % percent == 0):
                print(f'{i/percent}% complete.')
            with open('./temp/'+file, 'r') as fin:
                fout.write(fin.read())
    

    如果您追求这一点,还要尝试另一件事:改为以二进制模式打开文件('wb''rb')。这可能节省了文本模式字符解码/编码的无用层。我假设您只想将原始字节粘贴在一起。

    【讨论】:

      【解决方案3】:

      一种可扩展且可靠的方法是:

      • 配置 Amazon S3 存储桶以在新文件到达时触发 AWS Lambda 函数
      • 在 AWS Lambda 函数中,读取文件的内容并将其发送到 Amazon Kinesis Firehose 流。然后,删除输入文件。
      • 配置 Amazon Kinesis Firehose 流以缓冲输入数据并输出新文件,无论是基于时间段(最长 15 分钟)还是数据大小(最长 128MB)

      见:Amazon Kinesis Data Firehose Data Delivery - Amazon Kinesis Data Firehose

      不会每小时生成一个文件——文件的数量将取决于传入数据的大小。

      如果您需要创建每小时文件,您可以考虑对 Firehose 的输出文件使用 Amazon Athena。 Athena 允许您对存储在 Amazon S3 中的文件运行 SQL 查询。因此,如果输入文件包含日期列,它只能选择特定小时的数据。 (您可以为 Lambda 函数编写代码,为此添加一个日期列。)

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2016-08-22
        • 1970-01-01
        • 2019-12-19
        • 1970-01-01
        • 2020-03-16
        • 1970-01-01
        相关资源
        最近更新 更多