【发布时间】:2018-10-26 16:38:31
【问题描述】:
我正在使用 AWS Firehose 将传入记录转换为镶木地板。在一个示例中,我有 150k 条相同的记录进入 firehose,并且单个 30kb parquet 被写入 s3。由于 firehose 对数据的分区方式,我们在 parquet 中读取了一个辅助进程(由 s3 put 事件触发的 lambda)并根据事件本身的日期对其进行重新分区。在这个重新分区过程之后,30kb 的文件大小会跳到 900kb。
检查两个 parquet 文件-
- 元不变
- 数据没有变化
- 它们都使用 SNAPPY 压缩
- firehose parquet 由 parquet-mr 创建,pyarrow 生成的 parquet 由 parquet-cpp 创建
- pyarrow 生成的 parquet 有额外的 pandas 标头
完整的重新分区过程-
import pyarrow.parquet as pq
tmp_file = f'{TMP_DIR}/{rand_string()}'
s3_client.download_file(firehose_bucket, key, tmp_file)
pq_table = pq.read_table(tmp_file)
pq.write_to_dataset(
pq_table,
local_partitioned_dir,
partition_cols=['year', 'month', 'day', 'hour'],
use_deprecated_int96_timestamps=True
)
我想会有一些尺寸变化,但我惊讶地发现差异如此之大。鉴于我所描述的过程,什么会导致源拼花从 30kb 变为 900kb?
【问题讨论】:
-
如果没有可重复的例子,我们很难说出原因。我想不出有什么理由让我头脑发热
-
这可能与创建的文件数量有关。每个 Parquet 文件的固定开销为 4kb+。当您重新分区到太多文件时,这可能是来源之一。
标签: pandas parquet amazon-kinesis-firehose pyarrow