数据仓库优化手册
Hive性能优化小结
离线数仓任务优化
优化方向
- 尽早过滤无用数据
- 合理分配和利用资源,使每个task处理的数据量适中,避免启动过多或过少的task
- 减少shuffle或减少参与shuffle的数据量
- 避免数据倾斜,遇到倾斜的key进行打散处理
优化层面
业务层面
- 思考业务逻辑的合理性、可行性和必要性
- 计算量太大是否必须,是否可以减少参与计算的用户量或时间跨度
- 计算逻辑是否过于复杂,是否可以简化
模型层面
- 设计合理的数仓模型
- 是否有现成的数据可以使用或基于现成的数据进行加工
- 是否可以将整个计算逻辑进行合理拆分,降低每个子任务的复杂度,同时提高复用的可能性
- 维度退化,空间和时间的权衡
系统层面
- 遵循计算引擎建议的使用规则和参数设置
- 使用Spark3引擎,自动合并小文件
- 输入文件的存储格式、压缩格式、大小
- 输出文件的大小
- 启用压缩
- 分区、分桶
- 拉链表
- Yarn队列的设置
- 合适的计算引擎
- Task的内存设置
- Task处理的数据量
- Task的数量
- 并行度优化
- 调整参数减少Map数量
- 调整参数减少Reduce数量
SQL、代码层面
- 列裁剪,避免
SELECT * - 分区裁剪,使用分区字段过滤
- 条件限制
- 谓词下推
- Map端预聚合
- 大key的过滤
- 打散倾斜key
- 合适的Join方式
- 使用
DISTRIBUTE BY RAND控制分区中数据量 GROUP BY优化- 中间结果的缓存和复用
- 小文件优化
任务层面
- 减少任务依赖,尽可能缩短链路
- 业务链路/逻辑重构/改写
- 任务分级,任务数评估,错峰调度
- 任务依赖降级,周级别的任务依赖天级别,天级别依赖小时级别,小时级别依赖分钟级别
- 避免频繁创建任务
- 核心任务优先保证产出,双链路机制开启
- 耗时长的任务拆分成子任务,任务批次提交
- 资源动态扩容
- 资源腾挪调整
- 无用任务下线
Hive常用优化手段&参数
Spark常用优化手段&参数
自适应中Reduce参数控制
spark.sql.adaptive.shuffle.targetPostShuffleInputSize:控制任务Shuffle后的目标输入大小(以字节为单位)spark.sql.adaptive.minNumPostShufflePartitions:控制自适应执行中使用的Shuffle后最小分区数spark.sql.adaptive.maxNumPostShufflePartitions:控制Shuffle后分区的最大数量
合理设置单Partition读取数据量
SET spark.sql.files.maxPartitionBytes=xxxx;
合理设置Shuffle Partition的数量
SET spark.sql.shuffle.partitions=xxxx
使用COALESCE和REPARTITION调整Partition数量
sql
SELECT /*+ COALESCE(3) */ * FROM EMP_TABLE;
SELECT /*+ REPARTITION(3) */ * FROM EMP_TABLE;
SELECT /*+ REPARTITION(c) */ * FROM EMP_TABLE;
SELECT /*+ REPARTITION(3, dept_col) */ * FROM EMP_TABLE;
SELECT /*+ REPARTITION_BY_RANGE(dept_col) */ * FROM EMP_TABLE;
SELECT /*+ REPARTITION_BY_RANGE(3, dept_col) */ * FROM EMP_TABLE;
使用Broadcast Join
开启Adaptive Query Execution(Spark 3.0)
- 动态合并分区:Spark会根据分区的数据量将小数据量的多个分区合并成一个分区,提高资源利用率
spark.sql.adaptive.enabled:是否开启AQE优化spark.sql.adaptive.coalescePartitions.enabled:是否开启动态合并分区spark.sql.adaptive.coalescePartitions.initialPartitionNum:初始分区数spark.sql.adaptive.advisoryPartitionSizeInBytes:合并分区的推荐目标大小spark.sql.adaptive.coalescePartitions.minPartitionNum:合并之后的最小分区数
- 动态切换Join策略:在Join时,动态选择性能最高的Join策略,提高效率
spark.sql.adaptive.enabled:是否开启AQE优化spark.sql.adaptive.localShuffleReader.enabled:在不需要进行Shuffle重分区时,尝试使用本地Shuffle读取器
- 动态申请资源:当计算过程中资源不足时,自动申请资源
spark.sql.adaptive.enabled:是否开启AQE优化spark.dynamicAllocation.enabled:是否开启动态资源申请spark.dynamicAllocation.shuffleTracking.enabled:是否开启Shuffle状态跟踪
- 动态Join数据倾斜:Join时如果出现数据倾斜,动态调整分区的数据量,优化数据倾斜导致的性能问题
spark.sql.adaptive.enabled:是否开启AQE优化spark.sql.adaptive.skewJoin.skewedPartitionFactor:倾斜的膨胀系数spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes:倾斜的最低阈值spark.sql.adaptive.advisoryPartitionSizeInBytes:拆分粒度,以字节为单位
文件与分区
SET spark.sql.files.maxPartitionBytes=xxx:读取文件时一个分区接受多少数据spark.sql.files.openCostInBytes:文件打开的开销,小文件合并的阈值
CBO优化
spark.sql.cbo.enabled:是否开启CBO优化spark.sql.cbo.joinReorder.enabled:是否调整多表Join的顺序spark.sql.cbo.joinReorder.dp.threshold:设置多表Join的表数量的阈值
Hints优化
- Partitioning Hints Types:
COALESCE,REPARTITION,REPARTITION_BY_RANGE - Join Hints Types:
BROADCAST,MERGE,SHUFFLE_HASH,SHUFFLE_REPLICATE_NL
sql
SELECT /*+ COALESCE(3) */ * FROM t;
SELECT /*+ REPARTITION(3) */ * FROM t;
SELECT /*+ REPARTITION(c) */ * FROM t;
SELECT /*+ REPARTITION(3, c) */ * FROM t;
SELECT /*+ REPARTITION_BY_RANGE(c) */ * FROM t;
SELECT /*+ REPARTITION_BY_RANGE(3, c) */ * FROM t;
-- Join Hints for Broadcast Join
SELECT /*+ BROADCAST(t1) */ * FROM t1 INNER JOIN t2 ON t1.key = t2.key;
SELECT /*+ BROADCASTJOIN(t1) */ * FROM t1 LEFT JOIN t2 ON t1.key = t2.key;
SELECT /*+ MAPJOIN(t2) */ * FROM t1 RIGHT JOIN t2 ON t1.key = t2.key;
-- Join Hints for Shuffle Sort Merge Join
SELECT /*+ SHUFFLE_MERGE(t1) */ * FROM t1 INNER JOIN t2 ON t1.key = t2.key;
SELECT /*+ MERGEJOIN(t2) */ * FROM t1 INNER JOIN t2 ON t1.key = t2.key;
SELECT /*+ MERGE(t1) */ * FROM t1 INNER JOIN t2 ON t1.key = t2.key;
-- Join Hints for Shuffle Hash Join
SELECT /*+ SHUFFLE_HASH(t1) */ * FROM t1 INNER JOIN t2 ON t1.key = t2.key;
-- Join Hints for Shuffle-and-Replicate Nested Loop Join
SELECT /*+ SHUFFLE_REPLICATE_NL(t1) */ * FROM t1 INNER JOIN t2 ON t1.key = t2.key;
缓存表
- 使用
SQLContext.cacheTable(TableName)或DataFrame.cache缓存表 - 使用
SQLContext.uncacheTable(tableName)将表从缓存中移除 - 设置
spark.sql.inMemoryColumnarStorage.batchSize参数,默认10000,配置列存储单位
Group By优化
- 仅选择必要的字段进行
GROUP BY操作 - 尽可能将
GROUP BY字段类型保持一致 - 如果可能,可以将
GROUP BY字段进行哈希分区 - 如果使用字符串类型,可以考虑使用哈希函数减少字符串比较的开销
优化倾斜连接
- 数据偏斜会严重降低联接查询的性能
- 启用
spark.sql.adaptive.enabled和spark.sql.adaptive.skewJoin.enabled配置时,动态处理排序合并联接中的倾斜

实时数仓任务优化
Flink任务优化
时效优化
- 查看Kafka延迟监控:Flink消费上游的lag情况
- 分层和时延之间做好平衡和取舍,既保证复用性,又避免链路过长
- 乱序数据的处理
- 提前压测,应对流量高峰期,特别是大促场景下,提前做好资源保障、任务优化等措施
- 设置好延时基线,通过优化程序代码、资源、解决倾斜与反压等问题,使其控制在基线之内
- 指标监控,监控任务failover情况、checkpoint指标、GC情况、作业反压等,出现异常告警
- Flink链路延迟监控的LatencyMarker观测延迟情况
数据质量保障
数据一致性
- 正确性实时计算端到端的一致性,常用手段是通过输出幂等方式保障
- 离线与实时的一致性,需要保证使用数据源一致、加工业务逻辑一致
数据完整性
- 数据源层出现背压时,导致数据源头(MQ, Kafka)消息积压,积压严重时导致资源耗尽,进而导致数据丢失
- 数据处理层数据加工未按照需求进行加工,导致目标有效数据丢失
- 数据存储层的存储容量写满时,导致新数据无法继续写入导致数据丢失
- 数据加工正确性、数据加工及时性、数据快速恢复性构成数据完整性
数据加工正确性
- 数据源层原始数据包含不同联盟的点击数据,数据处理层过滤掉不需要的联盟点击数据,并将目标联盟的点击数据根据媒体和创意信息补齐当前点击所属的账号、计划、单元
- 业务层根据媒体、账号、计划、单元不同维度计算出对应的点击总量
数据加工及时性
- 目标源数据从产生到前端展示的时间需要控制在合理的时间范围内
数据快速恢复性
- 数据处理层因为消费程序性能问题导致消息积压,性能问题解决后数据挤压问题逐步得到缓解直到恢复正常水平
- 数据处理层因为消费程序bug导致程序崩溃,重启后数据消费正常
数据可监控性
- 数据流转路径中关键节点的关键状态可以有效监控
数据高可用性
- 数据不能因为灾难性的问题导致丢失造成不能使用的情况,因此需要考虑实时数据消费应用集群和存储集群的主备和可容灾
成本优化


计存成本优化
- CPU、内存、网络合理选择
- CPU进程绑定
- 数据分区分桶,选择合适的文件格式
- 数据分级与压缩存储
- 建立数据、数仓共享方案
- 计算引擎的选择
- 无用、低频数据下线或冷备
- 数据布局优化技术
- 优化系统
资源弹性伸缩容(云原生K8s等)
资源隔离,存算分离
稳定性保证

任务压测
- 提前压测应对流量高峰期,特别是大促场景下,提前做好资源保障、任务优化等措施
任务分级
- 制定保障等级,从任务影响面大小、数据使用方来划分,一般情况公司层面优先于部门层面,外部使用优先于内部使用,高优先级任务需要优先/及时响应、必要情况下做双链路保障机制
做好指标监控
- 指标监控,监控任务failover情况、checkpoint指标、GC情况、作业反压等,出现异常告警
高可用HA
- 整个实时Pipeline链路都应该选取高可用组件,确保理论上整体高可用;在数据关键链路上支持数据备份和重演机制;在业务关键链路上支持双跑融合机制
SLA保障
- 在确保集群和实时Pipeline高可用的前提下,支持动态扩容和数据处理流程自动漂移
弹性反脆弱
- 基于规则和算法的资源弹性伸缩;支持事件触发动作引擎的失效处理
监控预警
- 集群设施层面,物理管道层面,数据逻辑层面的多方面监控预警能力
自动运维
- 能够捕捉并存档缺失数据和处理异常,并具备定期自动重试机制修复问题数据
上游元数据变更抗性
- 上游业务库要求兼容性元数据变更;实时Pipeline处理显式字段
任务调度优化
- 减少任务依赖,尽可能缩短链路
- 业务链路/逻辑重构/改写
- 任务分级,任务数评估,错峰调度
- 任务依赖降级,周级别的任务依赖天级别,天级别依赖小时级别,小时级别依赖分钟级别
- 避免频繁创建小任务
- 核心任务优先保证产出,双链路机制开启
- 耗时长的任务拆分成子任务,任务批次提交
- 资源动态扩容
- 资源腾挪调整
- 无用任务下线
离线数仓任务延迟
1. 紧急修复故障
- 公司集群机器下线,尽快核查集群机器下线的原因,给出针对性解决方案
- 如果短时间内集群问题无法解决,尽快升级领导,讲清楚对业务的影响
2. 资源动态扩容
- 临时动态扩容是最基本的方法,如果这一步做不到,说明系统规划没做好
3. 资源腾挪调整
- 集群资源的使用存在2/8现象,将其他非核心任务的资源腾挪部分给核心任务
4. 任务分级、双链路切换
- 对所有任务的优先级进行排序,提升核心任务的调度优先级,调低其他任务的优先级
5. 任务代码优化
- 核心任务一般比较复杂,消耗的资源较多,有较多的优化空间
6. 降低任务依赖
- 通过剥离核心任务对数据仓库模型的依赖,为其量身定做一套数据处理逻辑,提升效率
7. 任务错峰调度、批量提交、无效任务下线,高消耗任务拆解
- 对于执行时间相差不大的任务,放在一起进行批量提交
- 对于低频或无效任务进行下线或降级
- 对于高消耗资源的任务,拆解成子任务进行提交
8. 核心任务容灾
- 采取异地异构的解决方案,确保核心任务万无一失
9. 做好解释工作
- 及时跟上级做好沟通汇报,协同各方做好故障的分析,给出可以落地的解决方案
- 做好核心任务延迟对业务造成的实际影响的评估,跟业务方做好解释工作
10. 转危机为机会
- 出故障对于数据仓库既是挑战,也是机会,可能会让公司重新审视数据仓库的价值

end
