尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

数据压缩在大数据中的核心价值:让Hive与Spark作业加速

数据压缩在大数据中的核心价值:让Hive与Spark作业加速 作为常年跟Hadoop、Spark、Hive打交道的数仓开发我越来越觉得“数据压缩”这词被大家理解得太窄了。很多人一听到压缩第一反应就是“省硬盘钱”。但在真正的大数据生产环境里压缩最核心的价值根本不是存储而是让处理速度飞起来。一个跑在Hive或Spark上的分析任务瓶颈往往不在CPU也不在算法而是卡在磁盘IO和网络传输上。把数据压缩了意味着从HDFS拉到计算节点、在shuffle阶段落盘、跨节点传输的数据量大幅减少整体作业时间可能直接砍半。这篇东西不聊虚的就是把我这些年在数仓、ETL、数据分析项目里用压缩加速的实操经验整理出来。内容包括压缩算法怎么选、Parquet/ORC这种列式存储和压缩怎么配合、Hive和Spark里具体怎么配置以及一堆只有踩过坑才知道的细节。适合正在做大数据开发、数据仓库搭建或者跑离线分析任务被IO折磨得头疼的同学参考。1. 先搞明白为什么压缩能让处理速度变快1.1 MapReduce和Spark的瓶颈根本是IO很多刚入行的同学有个误区觉得压缩是“压缩CPU换空间”跑任务时会变慢。这个想法在单机玩小数据时有点道理但在分布式计算框架里完全说反了。无论是MapReduce还是Spark核心计算模型都是“先把数据从磁盘读进内存处理后写到磁盘再通过网络传给下一阶段”。整个流程里数据移动的代价远比计算本身高。一次典型的Hive数仓查询磁盘读入、中间结果落盘、shuffle网络传输、最终结果写出IO占用可能占整个作业时间的70%以上。CPU在绝大多数时候是闲着的甚至要等数据从磁盘里搬进来。所以问题的本质是磁盘和网络的速度远远跟不上CPU和内存的处理速度。这时候如果能把数据体积压小那么同样的物理IO时间内我们就能读入更多“逻辑上”的数据。压缩消耗的那点CPU跟省下的IO时间比起来往往可以忽略不计。1.2 用压缩“换”IO这笔买卖什么时候划算我一般用一个粗略的判断方式压缩后的数据体积如果能减少50%以上在IO密集型作业里几乎稳赚即使只减少了30%在shuffle量很大的作业里也很划算。因为压缩省下的不只是HDFS存储还省了以下几笔账读数据阶段单次磁盘扫描能读到更多的有效数据shuffle阶段Map端输出压缩后写盘量和网络传输量都变小落盘阶段Spill到本地磁盘的中间数据更少结果写出如果最终结果也要落HDFS压缩后写得更快之前做过一个网约车订单数据的清洗任务原始数据量大概每天几十GB。不压缩跑一遍全量清洗要40多分钟在Hive表存储格式换成Parquet并开启Snappy压缩后同样逻辑的作业直接降到17分钟左右。这里面有列式存储的功劳也有压缩的功劳但两者是绑定的后面会细讲。注意如果集群的CPU已经持续90%以上内存又小压缩反而可能拖慢速度。这种情况要先扩容或者优化SQL不要盲目上压缩。1.3 存储成本只是这件事的副产品压缩后数据体积变小HDFS占用的块数变少这个当然也是好处。尤其是公司里的ODS层原始日志动辄上百TB不压缩根本存不起。但我觉得真正让压缩“有价值”的是它同时解决了存储成本和查询性能这两个问题。如果你的数仓建设还停留在文本文件加默认配置的阶段压缩是最容易入手、见效最快的一步。2. 大数据生态里的压缩算法到底怎么选2.1 主流的几个算法速览大数据框架里常见的压缩格式无非就这几个gzip、bzip2、lzo、snappy、zstdzstandard、lz4。它们各自特点很鲜明我按照实际项目里的使用频率逐个说。gzip压缩率比较高但是压缩和解压速度都偏慢。很多公司用gzip存离线冷数据因为不太关心查询速度只希望省空间。bzip2压缩率比gzip还高但速度更慢我基本只在归档场景见过它。lzo压缩和解压速度快而且支持分片splittable是早些年Hadoop生态里很流行的选择。有个大坑是它依赖libgpl需要自己装native库很多集群上默认没配好跑MapReduce时会报错。snappy在Google发布后迅速成为大数据领域的默认选项压缩速度极快解压速度也快压缩率比gzip差一些但非常均衡。Hadoop内置了对snappy的支持Hive、Spark、Parquet、ORC默认或者最推荐的压缩基本都是它。zstd这是近些年很值得关注的一个新选择。Facebook开源压缩率接近gzip但压缩速度远快于gzip而且在较高压缩级别下依然解压飞快。Spark 3.2、Hive 3.x系都已经很好支持了。lz4关注点就是极致的压缩和解压速度压缩率很低适合对速度敏感、对空间不敏感的场景。2.2 一张参数对比表看清权衡拿一个大概1GB的文本日志做测试集群环境不同会有浮动但趋势稳定我的经验数据大致如下算法压缩后大小压缩速度解压速度是否可分割适用场景gzip约250MB慢中否冷数据归档bzip2约200MB很慢很慢是极少用lzo约400MB快快需索引老集群历史方案snappy约350MB很快很快是容器格式内默认首选zstd约270MB中快很快是容器格式内新项目推荐lz4约600MB极快极快是实时计算这里要特别强调“是否可分割”。MapReduce/Spark读取文件时会把文件按HDFS块边界split成多个分片每个分片由不同Map或Task并行处理。如果一个压缩格式不可分割那么一个几百MB的压缩文件就只能被一个Task处理哪怕集群有几百个核心也只能眼睁睁看着一个核在跑其余全闲着。这就是压缩格式选错导致任务慢的直接原因。2.3 Hadoop、Spark、Hive里怎么配Hadoop层的全局压缩需要在core-site.xml配置property nameio.compression.codecs/name valueorg.apache.hadoop.io.compress.GzipCodec,org.apache.hadoop.io.compress.DefaultCodec,org.apache.hadoop.io.compress.SnappyCodec/value /property property nameio.compression.codec.snappy.native/name valuetrue/value /propertyMapReduce作业侧可以在mapred-site.xml或作业参数里开启Map输出压缩hadoop jar myjob.jar \ -Dmapreduce.map.output.compresstrue \ -Dmapreduce.map.output.compress.codecorg.apache.hadoop.io.compress.SnappyCodecSpark作业可以在提交时加参数spark-submit \ --conf spark.sql.parquet.compression.codecsnappy \ --conf spark.sql.orc.compression.codeczstd \ --conf spark.io.compression.codecsnappy \ --conf spark.shuffle.compresstrue \ --conf spark.shuffle.spill.compresstrue \这里要说一句Shuffle压缩往往比数据文件压缩更容易被忽略。很多同学只记得给表设压缩却忘了Spark默认shuffle是压缩的以及要确认shuffle压缩用的算法是否合适。shuffle产生的中间数据临时、量大、对速度要求高用snappy或lz4这类快速算法是正道不要用gzip去压shuffle数据。3. 存储格式和压缩是组合拳Parquet与ORC3.1 为什么列式存储比行式存储更适合压缩如果你还在用TextFile格式存Hive表那先别急着谈压缩算法选择因为存储格式带来的收益比压缩本身更明显。行式存储比如纯文本、SequenceFile把一条记录的各个字段连续写在一起。我们压整行数据时日志文本、数字、时间戳混在一起重复模式少压缩率自然一般。但列式存储如Parquet、ORC把每一列的数据连续存储相同类型的数据扎堆在一起。举个实际例子一个订单表的“订单状态”字段取值范围只有三五个用户ID字段可能会有大量重复前缀“下单时间”字段都是时间戳。这种重复度和规律性极高的列数据经过压缩算法能获得非常恐怖的空间收益通常比混合类型的行文本压缩率高20%到50%。列式存储带来的另一个关键优势是查询时的列裁剪Projection Pushdown。Hive/Spark读Parquet文件时可以直接跳过查询未涉及的列只读取需要的列块。这就意味着压缩节省的不只是“整体数据体积”而是真正读入IO的数据量大幅减少。比如一张表有20个字段查询只需要2个字段启用列裁剪后即使原表有1TB压缩后读取的数据可能只有几十MB。3.2 Parquet Snappy/Zstd 的标准操作目前数仓里最主流的方案就是Parquet加Snappy稳定、通用、社区支持好。我更推荐新集群或者允许调参的场景直接上Parquet加zstd压缩率更好跑批速度往往还能再快一点。Hive建表实操示例CREATE TABLE ods_order_daily ( order_id STRING, user_id STRING, driver_id STRING, city_id INT, order_time TIMESTAMP, order_amount DECIMAL(10,2), status STRING ) PARTITIONED BY (dt STRING) STORED AS PARQUET TBLPROPERTIES ( parquet.compression snappy );如果是Spark SQL也可以在SQL语句里设置SET spark.sql.parquet.compression.codecsnappy; INSERT OVERWRITE TABLE dwd_order_clean SELECT ... FROM ods_order_daily WHERE dt2024-06-01;还有一点容易被忽略写入端和读取端的压缩配置需要匹配。Parquet文件本身会在文件footer的元数据里记录压缩格式Hive和Spark一般都能自动识别并正确解压但当你用老版本Impala或者其他旧引擎去读高版本zstd压缩的Parquet时可能直接报Unsupported compression。生产环境里升级压缩算法前一定要确认下游查询引擎的版本都支持。3.3 ORC本身的索引机制和压缩能叠加如果Hive用得多ORC也是很值得考虑的选择。ORC自带轻量索引row group级别的min/max、布隆过滤在查询范围过滤时可以先读取索引判断这一块数据是否可能命中。这种“索引跳过”加上列式压缩在跑大表过滤统计时的性能提升是肉眼可见的。ORC建表的压缩配置示例CREATE TABLE dwd_user_action_orc ( user_id STRING, action_type STRING, action_ts BIGINT ) STORED AS ORC TBLPROPERTIES ( orc.compress zstd );注意ORC的历史版本里orc.compress支持NONE、ZLIB、SNAPPY、LZO高版本Hive才引入ZSTD。在Hive 2.x系的老环境里贸然设成zstd写入时不会报错但查询时会因为SerDe不认识直接失败。3.4 关于schema演进和压缩调整的一个提醒列式存储的压缩还有一个隐性好处修改表结构加列、删列的成本比文本格式低得多。因为列数据按列存储新增一列不会重写整个文件只新增独立的列块即可。对于ODS层经常要加字段的日常来说这是个省心且省IO的事。但对应的代价是——如果你频繁修改Parquet表的schema文件碎片会增加。所以建议常用表的压缩格式和schema在数仓建模阶段就定好上线后不要反复改格式数据量增长到一定规模后通过定期小文件合并compaction来保持文件块大小合适4. 实战Hive和Spark作业里配置压缩并验证效果4.1 一套可直接抄的数仓压缩方案我一般在数仓项目里按下述规范部署压缩策略这套结构在多个项目里验证过稳定性和性能都不错分层存储格式压缩算法用途ODS原始层Parquetgzip 或 zstd日志归档、明细存储空间优先DWD明细层Parquetsnappy 或 zstd日常ETL、查询分析速度与空间平衡DWS汇总层ORCzstd聚合结果、报表查询读取快ADS应用层Parquetlz4高频BI查询优先查询响应临时表/Shuffle不需要落地snappy/lz4中间过程数据IO速度优先这套方案的核心逻辑是越靠近底层数据量大且查询少用压缩率高的算法越靠近应用层查询越频繁用解压速度更快的算法。另外临时表不落HDFS就无所谓长期存储格式shuffle数据用快速算法即可。4.2 Hive跑批模拟与执行计划观察给一个生产里的实际场景。假设我们要清洗网约车订单数据原始表在ODS层存储为Parquetgzip清洗结果写入DWD层ParquetSnappySET hive.exec.compress.outputtrue; SET mapreduce.output.fileoutputformat.compress.codecorg.apache.hadoop.io.compress.SnappyCodec; SET mapreduce.map.output.compresstrue; SET mapreduce.map.output.compress.codecorg.apache.hadoop.io.compress.SnappyCodec; INSERT OVERWRITE TABLE dwd_order_clean PARTITION(dt2024-06-01) SELECT order_id, user_id, driver_id, city_id, order_time, order_amount, status FROM ods_order_daily WHERE dt2024-06-01 AND status IS NOT NULL;跑完后看两个东西一是作业的Counter里面Spilled records和Physical memory, 二是HDFS上新目录里的文件大小。替换压缩参数前后各跑一遍同样SQL对比记录表数据体积从多少降到多少Map阶段读取数据量Input下降了百分之多少作业总耗时缩短了多少是否有明显的GC或CPU打满现象用HDFS命令也可以直接看压缩效果和块分布hdfs dfs -ls /user/hive/warehouse/dwd_order_clean/dt2024-06-01 hdfs dfs -du -h /user/hive/warehouse/dwd_order_clean/dt2024-06-014.3 Spark作业里压缩参数的完整配置用Spark跑ETL是现在更流行的情况。我常用的Spark提交脚本示例spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.sql.parquet.compression.codeczstd \ --conf spark.sql.orc.compression.codeczstd \ --conf spark.io.compression.codecsnappy \ --conf spark.shuffle.compresstrue \ --conf spark.shuffle.spill.compresstrue \ --conf spark.rdd.compresstrue \ --conf spark.sql.shuffle.partitions400 \ job.py补充说明一下spark.rdd.compress。它控制RDD内存缓存数据在序列化之后是否压缩。在内存富余的时候可以关掉因为压缩会增加CPU开销但如果任务的内存比较紧张开这个参数能减少spill到磁盘的数据量算是用CPU换GC和IO的平滑方案。Spark作业跑的中间这些配置大家可以通过Spark UI观察Shuffle Read/Write的字节数变化。开启shuffle压缩后Shuffle Write一般能降一半以上。如果观察到的Shuffle Write流量没明显变化检查一下spark.shuffle.compress是否真的生效以及序列化是否走的是Kryo而不是Java序列化Kryo本身也会让数据更紧凑。4.4 结果数据的压缩率和耗时怎么看才标准对比测试压缩效果时注意控制变量相同集群资源同样的executor数量相同SQL逻辑相同输入数据分区多次运行取中位数避开其他作业干扰我自己测试时会把结果整理成格式场景数据体积读入数据量作业耗时说明Text 无压缩100GB100GB42min基线Parquet Snappy28GB4.2GB列裁剪后17min常用推荐Parquet ZSTD22GB4.2GB14min新推荐上面的“读入数据量”指的就是列裁剪后真正扫描到的数据量。这也是为什么列式存储和压缩组合起来效果远超单一措施的根本原因。5. 常见问题与排查技巧实录5.1 压缩文件不可分割导致的OOM和长尾之前接手过一个Hive作业输入文件是一批gzip压缩的日志每个文件100多MB。结果无论集群有多少资源这个作业始终只有对应文件数量的Task在跑而且每个Task都会把整个文件读入内存数据量大一点就OOM。原因就是gzip压缩文件本身不可分割。一个gzip文件必须从头读到尾才能完整解压所以一个文件只能由一个Map处理。如果文件大小是几百MB在HDFS上占了好几个块MapTask依然得读完整文件。这时候要么把大文件提前做合并和重写要么换用可分割的压缩容器格式Parquet/ORC天然解决了这个问题或者在压缩外层用SequenceFile包一下。经验总结文件型压缩gzip、lzo纯文件和容器内压缩Parquet、ORC的列块压缩是两个维度的东西。列式存储里的snappy压缩是列块级别的解压时不需要从文件开头顺序读取各Task可以分别开启各自的列块。这就是为什么“ParquetSnappy”很快而“纯Snappy文件”反而可能出现不可分割的坑。5.2 LZO native库缺失的“经典报错”Hadoop默认不带lzo的native库因为它的开源许可比较特殊。集群里如果没有提前装好hadoop-lzo依赖和对应的native二进制跑MapReduce直接报java.lang.RuntimeException: native-lzo library not available要解决也不是不行但每个节点都要装、配环境变量维护成本挺高。这也是我后来基本不推荐新项目用lzo的原因。实在要用就直接用Parquet/ORC里的lzo实现容器自身会处理依赖。5.3 snappy解压和gzip压缩混用为什么越跑越慢有个数仓同学把他的Hive中间表设置成输出gzip压缩Map输出也设成snappy结果作业比不压缩还慢。我看了一眼他的配置就明白了Map输出用snappy压缩快但Map阶段只是中间过程数据写盘确实少了但是Reduce端读取Map输出后要进行解压解压后数据量又膨胀回未压缩的大小这样洗数据时大尺寸的数据在落盘前才被gzip压缩gzip的压缩速度本来就慢压缩阶段耗时指数级上升最终整条链路的耗时完全被那个gzip压缩过程拖垮了正确的做法是中间路径全用快速算法snappy/lz4只有最终落地且长期存储的文件才用高压缩率算法gzip/zstd。比如ODS层写最终表可以用gzip但Map输出、shuffle中间结果一律snappy。5.4 磁盘空间不再紧张时也不要全换成高压缩率还有一个很容易犯的“好心办坏事”生产环境磁盘一度告急管理员把全链路压缩格式全部从snappy换成了gzip。存储问题确实解决了但所有下游报表查询、临时分析作业的耗时普遍上涨了30%以上业务部门直接来投诉。原因也很简单gzip解压速度比snappy慢得多查询类任务大量时间花在解压上而查询场景根本不在乎那点存储节省。后来我们按“冷热分离”的思路调整热表、高频查询表用snappy或lz4真正几个月前几乎不查的冷数据分区才统一改成gzip。两边兼顾。5.5 小文件问题会被压缩放大压缩可以减小体积但不会把小文件合并成大文件。如果生产管道里每天生成几千个几十KB的Parquet小块即使每个块都压缩Task数量依然爆炸NameNode内存压力依然大查询效率依然差。一个可靠的做法是跑定时任务把这些小文件合并成合理的块大小Parquet建议256MB到1GB之间。Hive可以通过INSERT OVERWRITE重写表触发合并Spark可以用coalesce或者repartition来控制输出文件数然后设置spark.sql.parquet.compression.codec让合并后的文件保持目标压缩格式。6. 一些让我少走弯路的经验细节6.1 千万不要对压缩数据再压缩一个很常见的场景日志数据在采集端已经用lz4压缩过一遍到了数仓后又用ParquetSnappy再压一遍。这样做的结果通常是CPU消耗翻倍存储和速度却没有进一步明显改善。因为第一次压缩已经把绝大部分冗余清掉了第二次压缩的收益非常低但CPU开销是实打实的。在建设数仓管道时评估一下上游数据源是否已做过压缩。如果源头已经是压缩格式优先考虑先解压再按目标格式存储或者直接跳过二次压缩而不是无脑统一压缩。这里说的“解压再压缩”确实会消耗IO但多数情况下比在压缩层上再叠压缩要划算。6.2 监控指标要盯着“有效数据吞吐量”看压测验证压缩优化效果时Spark UI里的“Input”指标不一定能直接反映有效性。在大数据任务里我更推荐关注以下指标Shuffle Write/Read字节数判断中间数据压缩是否生效Input Size / Records判断读入数据量是否因为列裁剪而下降任务处理时间中位数和长尾确认压缩格式是否影响并行处理GC时间压缩和解压会增加CPU但不应导致明显GC劣化CPU利用率峰值和均值确认压缩没有把CPU打满到新瓶颈如果看到任务“Input Size”没变但“作业耗时”下降那往往是解压速度快带来的收益如果“Input Size”明显下降那就是列裁剪和压缩率双重作用的结果。两种情形都说明优化有效但背后的调优方向不同。6.3 数据倾斜场景下压缩能帮你也不会帮你数据倾斜本身和压缩没有直接关系但压缩后数据体积变小会让倾斜问题更容易暴露。比如某些key的数据特别大在shuffle阶段压缩后一个Task依然要处理超大分区如果只看中间数据体积容易低估倾斜严重程度。遇到这种情况先看Spark UI里max task duration是否远大于中位数对于倾斜key可以考虑加盐salting或先做预聚合确保压缩和分区策略都合理后再谈加资源个人经验是压缩优化和大数据任务其他调优手段不是替代关系而是叠加关系。顺序应该是先看SQL和执行计划再调存储格式和压缩最后才是调整资源。6.4 zstd压缩级别怎么选zstd支持压缩级别参数从1到19还有0表示默认。我实测下来级别3到6在压缩率、压缩和解压速度的平衡上对大数仓场景是最合适的。级别8以上压缩率提升空间很小压缩耗时却可能翻倍。如果你用zstd处理历史冷数据分区可以开到9或12但处理日常流水和查询频繁的明细层用默认级别就好。在Parquet里设置zstd压缩级别可以在Spark里配置--conf spark.sql.parquet.compression.codec.zstd.level3在Hive里也可以通过parquet.compression.codec.zstd.level类似的tblproperties设置。注意并非所有版本支持使用前先确认引擎版本。6.5 依赖版本差异的排坑大数据框架版本对压缩的支持差异非常大。Hive 2.x和Spark 2.x默认用org.apache.hadoop.io.compress.SnappyCodec而Spark 3.x开始对zstd有原生而稳定的支持但需要确保spark-core自带对应的native库。如果你在Spark 2.x环境下想用zstd需要额外引依赖比如dependency groupIdcom.github.luben/groupId artifactIdzstd-jni/artifactId version1.5.2-2/version /dependency而且这个jar要分发到所有executor节点。我见过不少人在Spark 2.x里配置zstd后报NoClassDefFoundError就是因为没有正确分发依赖。6.6 压缩格式变更的迁移方案如果线上数仓已经跑了好几年想从Text/gzip迁到Parquet/Snappy不要期待一夜完成。稳妥的迁移步骤是先在ODS层新分区上采用新格式与老格式并存跑对比作业验证查询性能和新格式的稳定性逐步对历史分区执行INSERT OVERWRITE ... SELECT进行重写验证迁移后的数据质量和行数完全一致最后把老格式文件从HDFS上清理掉整个迁移过程不用停业务但一定要在业务低峰期执行并保留回滚计划。重写亿级分区耗时很长建议按时间分区逐个推进每个分区跑完先抽样校验。说实话做了这么多年大数据见过不少团队花大力气去调Spark执行参数、优化SQL却忘了最基础的数据压缩这一环。有时候把一个表从Text改成Parquet压缩从无到有查询速度提升的幅度远超那些精细调参。数据压缩是那种付出极少、回报极高的“便宜优化”但它需要从存储格式、压缩算法、查询场景几个维度通盘考虑。这也是我写这篇东西的核心原因形成一套适合自己的压缩策略你要踩的坑路我都标出来了。
返回列表