【发布时间】: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中。
我不知道如何提高执行时间,因为来自 parquet 的 INSERT 每个 partition_0 值需要 2 分钟。我单独尝试了select 以确保它不是INSERT 问题,并且需要1:50 分钟。所以问题是从S3读取数据。
如果我尝试为partition_0 选择一个不存在的值,则又需要大约 2 分钟,因此访问数据存在某种问题。不知道partition_0命名等是否算Hive分区格式。
编辑:添加 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