我们公司建设数据湖的过程中,遇到了严重的性能问题。

数据湖刚建好的时候,大家都很兴奋,终于有了统一的数据平台。但用了一段时间后,问题来了:查询慢、写入慢、资源消耗大。一个简单的查询,要跑十几分钟;一个ETL任务,要跑几个小时。数据团队天天被业务方催,压力很大。

我负责了这次性能优化,花了一个月时间,经过一系列调整,查询性能提升了10倍,写入性能提升了5倍。本文分享完整的优化过程。

一、数据湖背景

先说说我们数据湖的背景。

1. 建设目标

我们建设数据湖的目标:

  • 统一存储公司所有业务数据
  • 支持结构化、半结构化、非结构化数据
  • 支持多种计算引擎(Spark、Presto、Flink)
  • 支持批处理和流处理
  • 数据科学家和分析师可以自助查询

2. 技术栈

  • 存储:HDFS(后来迁移到对象存储S3)
  • 文件格式:Parquet
  • 计算引擎:Spark 3.0、Presto
  • 元数据:Hive Metastore
  • 数据格式:原始层(ODS)、明细层(DWD)、汇总层(DWS)、应用层(ADS)
  • 调度:Airflow

3. 数据规模

  • 总数据量:约500TB
  • 每天新增数据:约5TB
  • 表数量:约2000张
  • 每天任务数:约500个
  • 并发查询:约50个

4. 性能问题

主要的性能问题:

  • 查询慢:平均查询时间15分钟,最长的要1小时
  • 写入慢:ETL任务平均运行时间2小时,最长的要6小时
  • 资源消耗大:Spark任务经常OOM,CPU利用率低
  • 小文件多:HDFS上有几百万个小文件,NameNode压力大
  • 元数据慢:Hive Metastore查询慢,影响任务启动

二、性能分析

优化之前,先做全面的性能分析,找到瓶颈。

1. 慢查询分析

收集了最近一个月的慢查询,分析发现:

  • 大部分慢查询是全表扫描,没有分区裁剪
  • 有些查询扫描了大量不必要的列
  • Join操作没有广播,导致大量shuffle
  • 数据倾斜严重,某些task运行时间是其他的10倍
  • 文件格式不统一,有些是CSV,有些是JSON,有些是Parquet

2. 写入性能分析

分析ETL任务的性能:

  • 写入时产生大量小文件
  • 没有合理设置并行度,task数量过多或过少
  • 数据倾斜导致某些task写入慢
  • 压缩格式不合理,写入和读取都慢
  • 没有自动合并小文件的机制

3. 资源使用分析

分析集群资源使用:

  • CPU利用率低,平均只有30%
  • 内存利用率高,经常OOM
  • 磁盘IO高,尤其是shuffle阶段
  • 网络IO高,数据本地化差
  • 资源分配不合理,大任务和小任务抢资源

4. 元数据分析

分析Hive Metastore:

  • 表数量多,分区数量更多(有些表有上万个分区)
  • Metastore数据库查询慢
  • 没有定期清理过期分区
  • 统计信息不准确,导致优化器选错执行计划

三、优化过程

针对这些瓶颈,分步进行优化。

第一阶段:存储格式优化

1. 统一文件格式

问题:数据湖中文件格式不统一,有CSV、JSON、Parquet、ORC等。

CSV和JSON是文本格式,解析慢,不支持列裁剪和谓词下推,查询性能差。

优化:

  • 所有数据统一用Parquet格式
  • 原始数据入库时就转换成Parquet
  • 历史数据批量转换为Parquet
  • 禁止直接查询CSV和JSON文件

效果:查询性能提升了2-3倍,因为Parquet支持列裁剪和谓词下推,只读取需要的列和行。

2. 选择合适的压缩格式

问题:Parquet的压缩格式不统一,有些用gzip,有些用snappy,有些不压缩。

  • gzip:压缩率高,但解压慢,CPU消耗大
  • snappy:压缩率中等,解压快,CPU消耗小
  • 不压缩:读取快,但存储大,IO高

优化:

  • 冷数据(不常查询)用gzip,节省存储
  • 热数据(经常查询)用snappy,查询快
  • 统一用snappy作为默认压缩格式,平衡压缩率和速度

效果:查询性能提升了30%,存储空间减少了40%。

3. 合理设置文件大小

问题:小文件太多,HDFS上有几百万个小文件。

小文件的危害:

  • NameNode内存消耗大
  • 查询时打开文件的开销大
  • Map task数量多,调度开销大
  • 写入时产生大量小文件,影响性能

优化:

  • 目标文件大小设置为128MB-256MB
  • 写入时合理设置并行度,避免产生太多小文件
  • 定期合并小文件(用Spark的optimize命令,或者自己写合并任务)
  • 配置HDFS的小文件合并策略

效果:小文件数量减少了80%,查询性能提升了50%,NameNode压力大大减轻。

第二阶段:分区和分桶优化

1. 合理分区

问题:分区策略不合理。

  • 有些表没有分区,全表扫描
  • 有些表分区粒度过细(按天+小时+地区),分区数量太多
  • 有些表分区字段选择不合理,查询时无法分区裁剪

优化:

  • 所有大表都必须分区
  • 常用的查询条件作为分区字段(通常是日期)
  • 分区粒度适中:按天分区,不要按小时(除非数据量特别大)
  • 分区数量控制在合理范围(单表不超过1万个分区)
  • 定期清理过期分区

效果:大部分查询都能做分区裁剪,扫描的数据量减少了70%,查询性能提升了3倍。

2. 分桶优化

问题:Join操作慢,因为没有分桶,每次Join都要shuffle。

优化:

  • 对经常Join的大表,按Join key分桶
  • 分桶数量设置合理(通常是2的幂,如128、256)
  • Join的两张表用相同的分桶策略,可以避免shuffle
  • 用Bucketed Join代替普通Join

效果:Join性能提升了5倍,shuffle数据量减少了80%。

3. 数据布局优化

问题:数据在文件中的布局不合理,查询时需要读取大量不必要的数据。

优化:

  • 按查询频率高的列排序(Z-Ordering)
  • 对经常一起过滤的列,做多维聚簇
  • 用数据跳过(Data Skipping)技术,减少扫描的数据量
  • 定期优化数据布局(OPTIMIZE命令)

效果:查询扫描的数据量减少了60%,查询性能提升了2倍。

第三阶段:查询优化

1. 谓词下推和列裁剪

问题:很多查询没有利用谓词下推和列裁剪,扫描了大量不必要的数据。

优化:

  • 确保查询条件能下推到存储层
  • 只查询需要的列,不用SELECT *
  • 过滤条件尽量用分区字段和排序字段
  • 避免在过滤条件中对列做函数操作(如WHERE DATE(created_at) = '2022-09-01')

效果:查询扫描的数据量减少了50%,性能提升了1倍。

2. Join优化

问题:Join操作慢,主要是因为没有广播小表,数据倾斜。

优化:

  • 小表Join大表时,用Broadcast Join,把小表广播到所有节点
  • 合理设置广播阈值(如100MB以下的表自动广播)
  • 大表Join大表时,用Bucketed Join
  • 数据倾斜的Join,用Salting技术(给key加随机前缀)
  • 避免笛卡尔积Join

效果:Join性能提升了3-5倍。

3. 数据倾斜处理

问题:数据倾斜严重,某些task运行时间特别长,拖慢整个任务。

常见的倾斜场景:

  • Join key分布不均匀,某些key的数据量特别大
  • Group By key分布不均匀
  • 空值或默认值导致的倾斜

优化:

  • 倾斜的Join:用Salting,给key加随机前缀,分两阶段Join
  • 倾斜的Group By:先局部聚合,再全局聚合
  • 空值处理:空值单独处理,或者过滤掉
  • 开启Spark的倾斜处理参数(spark.sql.adaptive.enabled)

效果:任务运行时间稳定,不再有长尾task,整体性能提升了2倍。

4. AQE(自适应查询执行)

问题:Spark的静态优化器,不能根据运行时数据调整执行计划。

优化:

  • 开启AQE(Adaptive Query Execution)
  • 自动合并shuffle分区
  • 自动切换Join策略(SortMergeJoin转BroadcastJoin)
  • 自动优化数据倾斜
  • 合理设置AQE的参数

效果:查询性能平均提升了30%,很多任务不再需要手动调优。

第四阶段:写入优化

1. 合理设置并行度

问题:写入时并行度设置不合理。

  • 并行度太高:产生大量小文件
  • 并行度太低:写入慢,资源利用率低

优化:

  • 根据数据量动态调整并行度
  • 目标是每个task写入128MB-256MB的数据
  • 用repartition或coalesce控制输出文件数量
  • 开启动态分区写入时,注意小文件问题

效果:写入性能提升了2倍,小文件数量大大减少。

2. 批量写入和事务

问题:流式写入产生大量小文件,且没有事务保证。

优化:

  • 用Delta Lake或Iceberg,支持ACID事务
  • 流式写入时,合理设置checkpoint间隔
  • 定期合并小文件(OPTIMIZE)
  • 用MERGE INTO做upsert,避免重复数据

效果:写入性能提升了2倍,数据一致性得到保证。

3. 写入时压缩和排序

问题:写入时没有压缩和排序,导致文件大,查询慢。

优化:

  • 写入时就用snappy压缩
  • 写入时按查询频率高的列排序
  • 写入后自动生成统计信息
  • 写入后自动合并小文件

效果:写入后的文件更优,查询性能提升了50%。

第五阶段:资源和调度优化

1. 资源分配优化

问题:资源分配不合理,大任务和小任务抢资源。

优化:

  • 配置YARN的资源队列,不同类型的任务用不同队列
  • 大任务用大队列,小任务用小队列
  • 合理设置executor的内存和CPU
  • 开启动态资源分配,按需申请资源
  • 设置资源上限,防止单个任务占用所有资源

效果:资源利用率从30%提升到60%,任务等待时间减少了50%。

2. 内存优化

问题:Spark任务经常OOM,尤其是shuffle阶段。

优化:

  • 合理设置executor内存(4GB-8GB)
  • 增加shuffle内存比例
  • 开启off-heap内存,用于shuffle和缓存
  • 数据倾斜时,增加倾斜task的内存
  • 避免在driver端收集大量数据

效果:OOM减少了90%,任务稳定性大大提升。

3. 调度优化

问题:任务调度不合理,关键任务被阻塞。

优化:

  • 合理设置任务优先级
  • 关键任务用高优先级队列
  • 小任务和大任务分开调度
  • 错峰调度,避免高峰期资源竞争
  • 监控任务运行时间,及时发现异常

效果:关键任务的完成时间提前了30%,任务整体运行更稳定。

第六阶段:元数据优化

1. Metastore优化

问题:Hive Metastore慢,影响任务启动和查询。

优化:

  • Metastore数据库用高性能的MySQL或PostgreSQL
  • 给Metastore的表加索引
  • 定期清理过期的元数据
  • 限制单表的分区数量
  • 用Metastore的缓存,减少数据库查询

效果:Metastore查询速度提升了5倍,任务启动时间减少了80%。

2. 统计信息优化

问题:统计信息不准确,导致优化器选错执行计划。

优化:

  • 定期ANALYZE TABLE,更新统计信息
  • 收集列级统计信息(distinct count、null count、min/max)
  • 开启统计信息自动收集
  • 大表用采样统计,减少统计时间
  • 查询前检查统计信息是否最新

效果:执行计划更优,查询性能平均提升了20%。

四、优化效果

说说优化后的整体效果。

1. 查询性能

  • 平均查询时间:从15分钟降到1.5分钟,提升10倍
  • 最慢查询时间:从1小时降到10分钟,提升6倍
  • 并发查询能力:从50个提升到200个
  • 查询成功率:从90%提升到99%

2. 写入性能

  • ETL任务平均运行时间:从2小时降到25分钟,提升5倍
  • 数据延迟:从T+1降到小时级
  • 写入吞吐量:从每天5TB提升到每天20TB
  • 小文件数量:减少了80%

3. 资源使用

  • CPU利用率:从30%提升到60%
  • 内存OOM:减少了90%
  • 存储成本:减少了40%(压缩和小文件合并)
  • 集群规模:不需要扩容,节省了成本

4. 用户体验

  • 业务方满意度大幅提升
  • 数据分析师可以自助查询,不用等数据团队
  • 数据科学家可以快速迭代模型
  • 数据团队的运维压力大大减轻

五、踩过的坑

说说优化过程中踩过的坑。

坑一:分区粒度过细

有一张表,按天+小时+地区分区,分区数量超过了10万个。

结果:

  • Metastore压力大,查询分区慢
  • 每个分区的数据量很小,产生大量小文件
  • 查询时分区裁剪的开销很大

解决:

  • 改成按天分区,地区作为普通列
  • 分区数量控制在1万个以内
  • 历史分区合并

教训: 分区不是越细越好,要根据数据量和查询模式选择合适的粒度。


坑二:盲目使用Broadcast Join

开启自动广播后,有一次一张表统计信息不准,实际有2GB,但统计信息显示只有50MB,被自动广播了。

结果:

  • driver端OOM
  • 整个集群被拖慢

解决:

  • 调整广播阈值,从100MB降到50MB
  • 定期更新统计信息
  • 大表禁止广播,用hint强制SortMergeJoin

教训: 自动优化不是万能的,统计信息要准确,关键任务要手动验证执行计划。


坑三:合并小文件导致写入慢

开启自动合并小文件后,写入性能下降了。

原因:合并小文件需要额外的资源,和写入任务抢资源。

解决:

  • 合并小文件的任务和写入任务分开调度
  • 在低峰期合并小文件
  • 合理设置合并的触发条件

教训: 优化要权衡,不能为了查询性能牺牲写入性能。


坑四:数据倾斜处理不当

用Salting处理数据倾斜时,随机前缀的数量设置不合理。

  • 前缀太少:倾斜没有解决
  • 前缀太多:shuffle数据量增加,反而更慢

解决:

  • 根据倾斜key的数据量,动态调整前缀数量
  • 先分析数据分布,再设置参数
  • 用AQE自动处理倾斜,减少手动调优

教训: 数据倾斜处理要根据实际数据分布调整,不能一刀切。


坑五:迁移到对象存储后的性能问题

从HDFS迁移到S3后,查询性能下降了。

原因:

  • S3的IO延迟比HDFS高
  • S3不支持rename,写入时需要复制数据
  • 数据本地化差,网络IO高

解决:

  • 用S3的高性能存储类型
  • 增加文件大小,减少IO次数
  • 用缓存层(如Alluxio)加速热数据
  • 写入时用DirectOutputCommitter,避免rename

教训: 不同存储系统的特性不同,迁移后要重新调优。

六、数据湖性能优化的一般思路

总结一下数据湖性能优化的一般思路:

1. 存储层优化

  • 统一文件格式(Parquet或ORC)
  • 选择合适的压缩格式
  • 合理设置文件大小,避免小文件
  • 合理分区和分桶
  • 优化数据布局(排序、聚簇)

2. 计算层优化

  • 谓词下推和列裁剪
  • Join优化(Broadcast、Bucketed、Salting)
  • 数据倾斜处理
  • 开启AQE
  • 合理设置并行度

3. 资源层优化

  • 合理分配资源
  • 内存优化
  • 调度优化
  • 动态资源分配

4. 元数据优化

  • Metastore性能优化
  • 统计信息收集
  • 分区管理
  • 元数据缓存

5. 持续监控和调优

  • 监控慢查询和慢任务
  • 定期分析性能瓶颈
  • 持续优化,不是一劳永逸
  • 建立性能基线,对比优化效果

七、写在最后

这次数据湖的性能优化,把查询性能提升了10倍,写入性能提升了5倍,效果很明显。

但性能优化不是一次性的工作,而是持续的过程。数据在增长,业务在变化,新的性能问题会不断出现。我们需要建立完善的监控体系,持续关注性能,持续优化。

数据湖的性能优化,涉及存储、计算、资源、元数据等多个层面。任何一个层面的问题,都可能影响整体性能。需要全面分析,找到真正的瓶颈,针对性优化。

2022年了,数据湖已经成为很多公司的数据基础设施。但建好数据湖只是第一步,用好数据湖才是关键。性能优化,是用好数据湖的重要保障。

最后,用一句话总结:"数据湖性能优化,存储是基础,计算是关键,资源是保障,元数据是灵魂。全面分析,持续优化,才能让数据湖又快又稳。"

愿你的数据湖,查询如风,写入如电。