Spark与Kinesis Firehose生成Parquet文件大小差异原因排查
咱们来一步步分析你遇到的Spark与Firehose生成Parquet文件体积差异的问题,从你提供的元数据和文件细节来看,核心原因是单行数组成员密度过高带来的编码开销,再加上一些额外的元数据和编码细节差异,具体如下:
1. 行数据分布是核心差异点
先看两个文件的行组数据:
- Firehose生成的文件有156行,每行平均只包含几百个数组元素(比如
dfpts列算上空值总元素数约16万,平均每行1000出头) - Spark生成的文件只有11行,但每行平均塞了数千个数组元素(
dfpts列总元素数6万多,平均每行接近6000个)
Parquet存储数组这类重复字段时,靠**重复级别(Repetition Level)**来标记每个元素属于哪一行。如果单行数组成员太少,RLE游程编码可以高效压缩重复的行标记;但如果每行塞了几千个元素,重复级别需要不断标记“这个元素还是属于当前行”,RLE的压缩效率会暴跌,直接导致存储开销飙升。
2. Spark专属元数据的额外开销
Spark生成的Parquet文件里会自带一段org.apache.spark.sql.parquet.row.metadata的JSON元数据,用来记录Spark的Schema信息。单看这个文件的话这段体积不大,但如果是大量小Parquet文件累积,这部分开销会被放大。
至于Schema里的命名差异(list.element vs bag.array_element),这个只是Spark和Hive对数组结构的命名规范不同,本质都是Parquet的重复组实现,不会直接导致体积差异。
3. 空值处理与编码效率的区别
Firehose的数组列有不少空值(比如udids列有58个空值),Parquet对空值是用bitpacking标记的,空值本身不占元素存储的空间;而Spark的数组列没有空值,所有元素都得完整存储。
另外,虽然两者都用了PLAIN_DICTIONARY编码,但Firehose的总元素数更多,相同元素重复出现的概率更高,字典复用率也就更好,压缩效率自然更高。
4. Parquet版本的细微优化差异
Spark用的是parquet-mr 1.8.3,Firehose是1.8.1,这两个小版本在字典编码、重复级别编码的细节上可能有优化差异,也会影响最终的存储体积。
给你的优化建议
如果想缩小Spark生成的Parquet文件体积,可以试试这几个方向:
- 控制每个Parquet文件的行数:通过
spark.sql.files.maxRecordsPerFile参数,让每个文件的行数更接近Firehose的水平,降低单行数组成员的密度,减少重复级别的编码开销。 - 确认字典编码配置:Spark默认开启了
spark.sql.parquet.enableDictionary,可以确认这个配置没被修改,确保字典编码能充分发挥作用。 - 合并小文件:如果Spark生成了大量小Parquet文件,元数据的累积开销会很可观,用
coalesce或repartition合并文件,提升整体的压缩效率。 - 检查数据写入逻辑:看看是不是Spark任务里把多条数据合并成了单行数组成员,尽量让行分布更均匀,避免单行数组成员过载。
内容的提问来源于stack exchange,提问作者jph

