【问题标题】:Parquet data to AWS Redshift slowParquet 数据到 AWS Redshift 速度很慢
【发布时间】:2021-03-23 05:11:06
【问题描述】:

我想将 S3 parquet 文件中的数据插入 Redshift。

parquet 中的文件来自读取 JSON 文件、将它们展平并存储为 parquet 的进程。为此,我们使用pandas dataframes

为此,我尝试了两种不同的方法。第一个:

COPY schema.table
FROM 's3://parquet/provider/A/2020/11/10/11/'
IAM_ROLE 'arn:aws:iam::XXXX'
FORMAT AS PARQUET;

它返回:

Invalid operation: Spectrum Scan Error
error:  Spectrum Scan Error
code:      15001
context:   Unmatched number of columns between table and file. Table columns: 54, Data columns: 41

我了解这个错误,但我没有一个简单的选项来修复它。 如果我们必须从 2 个月前重新加载,该文件将只有 40 列,因为在给定的数据上,我们只需要这些数据,但表已经增加到 50 列。 所以我们需要一些自动的东西,或者我们至少可以指定列。

然后我应用了另一个选项,即使用AWS Redshift Spectrum 执行SELECT。我们知道该表有多少列使用系统表,我们现在将文件的结构再次加载到Pandas dataframe。然后我可以将两者结合起来,使其具有相同的结构并进行插入。

它工作正常,但速度很慢。

选择看起来像:

SELECT fields
FROM schema.table
WHERE partition_0 = 'A'
  AND partition_1 = '2020'
  AND partition_2 = '11'
  AND partition_3 = '10'
  AND partition_4 = '11';

当我使用以下方法检查时,已经添加了分区:

select *
from SVV_EXTERNAL_PARTITIONS
where tablename = 'table'
  and schemaname = 'schema'
  and values = '["A","2020","11","10","11"]'
limit 1;

我每小时大约有 170 个文件,包括 json 和 parquet 文件。该进程列出S3 json path中的所有文件,并将它们处理并存储在S3 parquet path中。

我不知道如何提高执行时间,因为来自 parquetINSERT 每个 partition_0 值需要 2 分钟。我单独尝试了select 以确保它不是INSERT 问题,并且需要1:50 分钟。所以问题是从S3读取数据。

如果我尝试为partition_0 选择一个不存在的值,则又需要大约 2 分钟,因此访问数据存在某种问题。不知道partition_0命名等是否算Hive分区格式。

编辑: AWS Glue 爬虫表规范

编辑:添加 SVL_S3QUERY_SUMMARY 结果

step:1
starttime: 2020-12-13 07:13:16.267437
endtime: 2020-12-13 07:13:19.644975
elapsed: 3377538
aborted: 0
external_table_name: S3 Scan schema_table
file_format: Parquet         
is_partitioned: t
is_rrscan: f
is_nested: f
s3_scanned_rows: 1132
s3_scanned_bytes: 4131968
s3query_returned_rows: 1132
s3query_returned_bytes: 346923
files: 169
files_max: 34
files_avg: 28
splits: 169
splits_max: 34
splits_avg: 28
total_split_size: 3181587
max_split_size: 30811
avg_split_size: 18825
total_retries:0
max_retries:0
max_request_duration: 360496
avg_request_duration: 172371
max_request_parallelism: 10
avg_request_parallelism: 8.4
total_slowdown_count: 0
max_slowdown_count: 0

添加查询检查

查询:37005074(使用 pycharm 在 localhost 中选择) 查询:37005081(在 AIRFLOW AWS ECS 服务中插入)

STL_QUERY显示两个查询都需要大约 2 分钟

select * from STL_QUERY where query=37005081 OR query=37005074 order by query asc;

Query: 37005074 2020-12-14 07:44:57.164336,2020-12-14 07:46:36.094645,0,0,24
Query: 37005081 2020-12-14 07:45:04.551428,2020-12-14 07:46:44.834257,0,0,3

STL_WLM_QUERY显示没有排队时间,全部在执行时间

select * from STL_WLM_QUERY where query=37005081 OR query=37005074;

Query: 37005074 Queue time 0 Exec time: 98924036 est_peak_mem:0
Query: 37005081 Queue time 0 Exec time: 100279214 est_peak_mem:2097152

SVL_S3QUERY_SUMMARY表明查询在 s3 中需要 3-4 秒

select * from SVL_S3QUERY_SUMMARY where query=37005081 OR query=37005074 order by endtime desc;

Query: 37005074 2020-12-14 07:46:33.179352,2020-12-14 07:46:36.091295
Query: 37005081 2020-12-14 07:46:41.869487,2020-12-14 07:46:44.807106

stl_return比较每个查询的最小开始和最大结束。正如SVL_S3QUERY_SUMMARY 所说,3-4 秒

select * from stl_return where query=37005081 OR query=37005074 order by query asc;

Query:37005074  2020-12-14 07:46:33.175320 2020-12-14 07:46:36.091295
Query:37005081  2020-12-14 07:46:44.817680 2020-12-14 07:46:44.832649

我不明白为什么 SVL_S3QUERY_SUMMARY 显示运行频谱查询只需 3-4 秒,但随后 STL_WLM_QUERY 说执行时间约为 2 分钟,正如我在本地主机和生产环境中看到的那样......改进它,因为stl_return 表明查询返回的数据很少。

解释

XN Partition Loop  (cost=0.00..400000022.50 rows=10000000000 width=19608)
  ->  XN Seq Scan PartitionInfo of parquet.table  (cost=0.00..22.50 rows=1 width=0)
        Filter: (((partition_0)::text = 'A'::text) AND ((partition_1)::text = '2020'::text) AND ((partition_2)::text = '12'::text) AND ((partition_3)::text = '10'::text) AND ((partition_4)::text = '12'::text))
  ->  XN S3 Query Scan parquet  (cost=0.00..200000000.00 rows=10000000000 width=19608)
"        ->  S3 Seq Scan parquet.table location:""s3://parquet"" format:PARQUET  (cost=0.00..100000000.00 rows=10000000000 width=19608)"

svl_query_report

select * from svl_query_report where query=37005074 order by segment, step, elapsed_time, rows;

【问题讨论】:

  • 为什么 2 分钟对你来说太慢了(你不能一夜之间加载它们)?你有多少个文件?每个包含多少行?每个分区只有一个文件吗?是什么阻止它一次加载多个文件——是因为历史文件有不同的格式吗?它会
  • @JohnRotenstein 2 分钟很慢,因为该过程是每小时一次。因此,要加载 30 个不同的 partitions_0 值需要一小时......而且 partitions_0 的数量仍在增加。行依赖,它来自 s3 中的 json 文件,这些文件是通过 kinesis 创建的。所以每个文件的数据大约是 30kb。我们设置了 1mb 或 1 分钟来限制 kinesis 的大小。我们加载每个文件,因为我们列出了 s3 路径中的所有文件,所以自然迭代是按文件执行,因为我们有一个列表。无论如何,我想尝试的是将所有文件加入一个 df 中。
  • 为什么要加载 30 个不同的分区?如果数据在一小时内生成,我建议在一次 COPY 操作中加载所有文件。这可以通过将 Redshift 指向包含文件的目录来完成,它还将加载子目录中的所有文件。
  • 我加载了 30 个不同的分区,因为每个分区都是一个提供者,每个分区都有自己的表。复制操作不起作用。正如我在帖子中所说:如果我们必须从 2 个月前重新加载,该文件将只有 40 列,因为在给定的数据上,我们只需要这些数据,但表已经增加到 50 列。所以我们需要一些自动的东西,或者我们至少可以指定列。

标签: python-3.x amazon-web-services amazon-s3 amazon-redshift parquet


【解决方案1】:

就像在您的其他问题中一样,您需要更改对象上的键路径。仅在 keypath 中有“A”是不够的 - 它需要是“partition_0=A”。这就是 Spectrum 知道对象是否在分区中的方式。

您还需要确保您的对象大小合理,否则如果您需要扫描其中的许多对象,则速度会很慢。打开每个对象都需要时间,如果您有许多小对象,则打开它们的时间可能比扫描它们的时间长。只有当您需要扫描许多文件时,这才是一个问题。

【讨论】:

  • 我将更改键路径。只是改名字,还是需要做其他操作?关于文件我已经做了一个没有用的测试,但可能是由于密钥路径。所以我会先尝试添加keypath,看看它是如何改进的。
  • 我对其中一个路径进行了更改,但仍然需要将近 2 分钟(尚未尝试创建更大的文件)。难道只是改变一条路径并没有加快速度吗?我刚改了:alter table schema.table add partition (partition_0 = 'A', partition_1 = '2020', partition_2 = '11', partition_3 = '10', partition_4 = '11') location 's3://parquet/ provider_parquet/partition_0=A/partition_1=2020/partition_2=11/partition_3=10/partition_4=11​​/';并将文件从旧移动到这个
  • 您能否发布您的“CREATE EXTERNAL TABLE”语句以及任何分区修饰符?
  • 宾果游戏。请记住,S3 中没有“文件夹”——它是平面存储。桶和钥匙。现在,当键中包含“/”时,读取存储桶的内容将其视为方便访问的分隔符 - 它不是数据在最低级别存储方式的一部分。查找 1 个键很快,因为散列值马上就知道了。但是搜索文件夹很慢(在您的情况下),因为 S3 必须列出整个存储桶并将密钥与您想要的密钥模式进行比较。您要修复的路径是 1)更少的小对象和 2)在存储桶之间拆分数据。你可能两者都需要。
  • 在我的记忆中单击的另一种攻击方法 - 您可以在定义外部表时使用清单文件而不是“文件夹”。这应该意味着 Redshift 拥有表中所有对象的所有文件键,而不必列出所有 1100 万个键。您需要将清单文件与配置分区的 alter table 语句一起使用。这是新的,迄今为止我还没有尝试过使用大桶的解决方案,所以如果这解决了这个问题,请通知大家。我相信外部表上的清单旨在解决这个问题(以及其他一些问题)。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-04-12
  • 1970-01-01
  • 2017-06-03
  • 2016-06-28
  • 2015-03-26
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多