【问题标题】:Alternative to create more than 100 partitions on Athena CTAS在 Athena CTAS 上创建 100 多个分区的替代方法
【发布时间】:2020-02-21 08:07:12
【问题描述】:

我目前正在根据存储在 Amazon S3 中的信息创建一些新表。第一次使用 AWS,今天我了解到 Amazon Athena 无法通过 CTAS 查询创建超过 100 个分区。

我正在使用 sql 进行转换,它运行良好,但需要一种方法来一次存储 100 多个分区以使过程更可靠。

我将分区设置为日期,因此如果我需要重新创建表以通过 sql 加载大量数据(我有转换),那么在 4 个月内我的过程将失败。

知道如何实现这一目标吗?

【问题讨论】:

  • 你特别需要按天分区吗?你可以按月分区吗?您的转换是在一天内执行的吗?
  • 最终表上的查询通常按天过滤,并且按天而不是按周或按月维护似乎更容易删除和重新生成,但我愿意接受建议。谢谢!
  • 只是历史数据,还是应该每天/每月更新一次?
  • @IlyaKisil 它每天/每天更新两次

标签: amazon-web-services amazon-s3 amazon-athena


【解决方案1】:

最好的选择是为此任务编写一个 Glue ETL (spark) 作业并使用 spark sql 执行所需的转换。这样您仍然可以使用现有的 sql 查询。

然后您可以将处理后的输出写回某个 S3 路径。 Spark 允许您创建任意数量的分区。它还允许将新处理的数据附加到已处理的数据中,从而允许您仅加载和转换新数据。

ETL 完成后,创建一个指向上述 S3 路径和所需分区的外部表。这将是一个时间步骤(创建外部表)。您只需在每次粘合作业后更新此外部表中的分区信息。

总而言之,您需要执行以下操作:

  • 创建要在 Glue ETL 上执行的 spark 脚本,该脚本将读取每日源数据,应用所需的转换并将处理后的数据写入 S3 上的新分区中。这个脚本可以很容易地被制作成接受日期作为输入,并且是一次性活动。

  • 创建一个指向 S3 上已处理数据的外部表。这也是一次性活动。

  • 在每次 Glue ETL 作业后对上述外部表执行 MSCK 修复命令以更新新分区。

参考资料:

AWS Glue ETL documentation

AWS Athena - Create external table

AWS Athena - Update partiotion

【讨论】:

  • 谢谢@harsh-bafna,很抱歉延迟接受答案,我将使用 Glue,你在那儿提出了一个很好的观点,我有一个小问题,我可以访问一个问题像我在雅典娜做的那样的数组/结构? client.id 或 element[1].name 我现在有一个错误,我很难理解 spark 消息的页面和页面。如果找不到解决方案,会再发一个帖子
  • @Alejandro : 可能这个线程可以提供帮助:stackoverflow.com/questions/34069282/…
【解决方案2】:

假设您想使用 CTAS 查询处理 4 个月的数据,但您需要按天对其进行分区。如果您在单个 CTAS 查询中执行此操作,您最终将得到大约 4 x 30 = 120 个分区,因此,由于AWS limitations,查询将失败,正如您所提到的。

相反,您可以一次处理每个月的数据,这样可以保证您一次拥有少于 31 个分区。但是,每个 CTAS 查询的结果应该在 S3 上具有唯一的外部位置,即,如果您想在 s3://bukcet-name/data-root 下存储多个 CTAS 查询的结果,您需要在 external_location 下为 WITH 下的每个查询扩展此路径条款。对于您的案例,显而易见的选择是完整日期,例如:

s3://bukcet-name/data-root
├──2019-01            <-- external_location='s3://bukcet-name/data-root/2019-01'
|   └── month=01
|       ├── day=01
|       |   ...
|       └── day=31
├──2019-02            <-- external_location='s3://bukcet-name/data-root/2019-02'
|   └── month=02
|       ├── day=01
|       |   ...
|       └── day=28
...

但是,现在您得到了 4 个不同的表。因此,您要么需要查询不同的表,要么必须进行一些后处理。本质上,您将有两个选择

  1. 将所有新文件移动到带有AWS CLI high-level commands 的公共位置,后面应该跟着MSCK REPAIR TABLE,因为输出“目录”结构遵循HIVE 分区命名约定。例如来自

    s3://bukcet-name/data-staging-area
    ├──2019-01        <-- external_location='s3://bukcet-name/data-staging-area/2019-01'
    |   └── month=01
    |       ├── day=01
    |       |   ...
    

    你会复制到

    s3://bukcet-name/data-root
    ├── month=01
    |  ├── day=01
    |  |   ...
    |  └── day=31
    ├── month=02
    |   ├── day=01
    |   |   ...
    |   └── day=28
    
  2. 使用 AWS Glue 数据目录进行操作。这有点棘手,但主要思想是您定义一个根表,其位置指向s3://bukcet-name/data-root。然后在执行 CTAS 查询后,您需要将有关分区的元信息从创建的“暂存”表复制到根表中。此步骤将基于AWS Glue API,例如用于 Python 的boto3 库。特别是,您将使用get_partitions()batch_create_partition() 方法。

无论您选择哪种方法,您都需要使用某种作业调度软件,尤其是因为您的数据不仅仅是历史数据。我建议为此使用Apache Airflow。它可以看作是 Lambda 和 Step Functions 组合的替代方案,它是完全免费的。有很多博客文章和文档可以帮助您入门。例如:

  • Medium post:使用 Airflow 自动执行 AWS Athena 查询并在 S3 周围移动结果。
  • Airflow 安装完整指南,link 1link 2

您甚至可以设置integration with Slack 在查询以成功或失败状态终止时发送通知。

注意事项:

通常,您无法明确控制 CTAS 查询将创建多少文件,因为 Athena 是一个分布式系统。另一方面,您不想拥有很多小文件。所以可以尝试使用"this workaround",它在WITH 子句中使用bucketed_bybucket_count 字段

CREATE TABLE new_table
WITH (
    ...
    bucketed_by=ARRAY['some_column_from_select'],
    bucket_count=1
) AS (
    -- Here goes your normal query 
    SELECT 
        *
    FROM 
        old_table;
)

或者,减少分区数量,即在月级别停止。

【讨论】:

  • 听起来很有趣,我试图避免使用更多的外部工具,因为我正在添加更多的故障点。现在我有用于提取的 NIFI 并在 Athena 中进行转换,监控管道是否有错误变得更加棘手。谢谢!
  • 我明白)Airflow 提供了很好的 UI 用于监控和重置失败的任务,这是一个巨大的好处。它甚至有一个dedicated operator 用于向 Athena 发送查询并在 S3 上进行操作)。当我的 ETL 作业开始时,我开始使用它,不得不监控)
猜你喜欢
  • 1970-01-01
  • 2021-12-04
  • 2019-05-03
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-12-09
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多