数仓架构师必备技能点下篇
大数据技术提升指南
如果你想提升大数据技术能力,寻找专业指导,持续学习,并找到一份满意的工作,欢迎加入我的知识星球。我已经开通知识星球大半年,星球上沉淀了大量有价值的内容。
星球的定位是专注于大数据内容的干货输出,帮助大家提升大数据技术,带领大家共同学习。星球内容及语雀文档都是经过我精心总结和整理的。



架构迭代
Kubernetes (k8s)
优点
- 统一运维:公司统一化运维,有专门的部门运维 K8S。
- CPU 隔离和安全性:K8S Pod 之间 CPU 隔离,实时任务不相互影响,更加稳定。
- 存储计算分离:Flink 计算资源和状态存储分离,计算资源能够和其他组件资源进行混部,提升机器使用率。
- 弹性扩缩容:大促期间能够弹性扩缩容,更好地节省人力和物力成本。
缺点
- 学习成本高、运维成本高:需要专业的运维。
- 与大数据生态集成相对欠缺。
开源大数据产品
- CloudEon
- KubeBlocks
Spark on K8S Native 和 Operator 的区别

Flink on K8S Native 和 Spark on Operator 的区别

湖仓一体
数据湖
- 数据湖解决了离线数仓什么问题?
- 数据湖能够解决离线数仓所面临的一些问题,包括但不限于以下几点:
- 数据类型和格式的限制:离线数仓通常只能处理结构化数据,而无法承载其他类型的数据。
- 数据存储和处理效率:离线数仓需要批量处理大量历史数据,因此需要在处理过程中进行数据移动和转换,导致处理效率低下。
- 数据质量:离线数仓在数据处理过程中容易发生数据质量问题,如数据重复、数据丢失、数据错误等。
- 数据湖读时模式:在读取数据时才检验数据。
- 增量更新问题。
- 数据湖能够解决离线数仓所面临的一些问题,包括但不限于以下几点:
相比之下,数据湖能够解决这些问题,因为它具有以下特点:
-
多样化的数据类型和格式:数据湖能够承载各种类型和格式的数据,包括结构化数据、半结构化数据和非结构化数据。
-
高容量和弹性的存储:数据湖在存储方面比较灵活,可以存储大量的数据,并且能够随时扩容和缩容。
-
处理速度和效率:数据湖具有在大规模数据处理方面的优势,因为它可以利用云计算和分布式技术来实现数据的处理和计算,从而提高处理速度和效率。
-
数据质量保证:数据湖中的数据质量能够得到保证,因为它具有数据治理的功能,可以对数据进行验证、清洗、补充和校准。
-
一款好的数据湖需要具备哪些功能?
-
- 能够存储任意数据,结构化、非结构化、半结构化的数据。
-
- 数据具有可访问性和开放性。
-
- 丰富的计算引擎,从批处理、流式计算、交互式分析到机器学习,各类计算引擎都属于数据湖应该囊括的范畴。
-
- 多模态的存储引擎。理论上,数据湖本身应该内置多模态的存储引擎,以满足不同的应用对于数据访问需求(综合考虑响应时间/并发/访问频次/成本等因素)。但是,在实际的使用过程中,数据湖中的数据通常并不会被高频次的访问,而且相关的应用也多在进行探索式的数据应用,为了达到可接受的性价比,数据湖建设通常会选择相对便宜的存储引擎(如S3/OSS/HDFS/OBS),并且在需要时与外置存储引擎协同工作,满足多样化的应用需求。
-
- 支持记录级别的 update/delete,以增量更新数据。
-
- 支持并发读写和事务的 ACID 特性;支持 MVCC。
-
- 支持历史版本回溯。
-
- 支持模式约束和演化 schema evolution, schema enforcement。
-
- 灵活的元数据管理和组织形式。
-
-
数据中台、数据仓库、大数据平台、数据湖的关键区别是什么?
-
- 基础能力上的区别
- 数据平台:提供的是计算和存储能力。
- 数据仓库:利用数据平台提供的计算和存储能力,在一套方法论的指导下建设的一整套的数据表。
- 数据中台:包含了数据平台和数据仓库的所有内容,将其打包,并且以更加整合以及更加产品化的方式对外提供服务和价值。
- 数据湖:一个存储企业各种各样原始数据的大型仓库,包括结构化和非结构化数据,其中湖里的数据可供存取、处理、分析和传输。
-
- 业务能力上的区别
- 数据平台:为业务提供数据主要方式是提供数据集。
- 数据仓库:相对具体的功能概念是存储和管理一个或多个主题数据的集合,为业务提供服务的方式主要是分析报表。
- 数据中台:企业级的逻辑概念,体现企业数据产生价值的能力,为业务提供服务的主要方式是数据 API。
- 数据湖:数据仓库的数据来源。
-
总的来说,数据中台距离业务更近,数据复用能力更强,能为业务提供速度更快的服务,数据中台在数据仓库和数据平台的基础上,将数据生产为一个个数据 API 服务,以更高效的方式提供给业务。数据中台可以建立在数据仓库和数据平台之上,是加速企业从数据到业务价值的过程的中间层。
-
数据湖、数据仓库、湖仓一体之间的区别?
- 湖仓一体是一种新型开放式架构,将数据湖和数据仓库的优势充分结合,它构建在数据湖低成本的数据存储架构之上,又继承了数据仓库的数据处理和管理功能,打通数据湖和数据仓库两套体系,让数据和计算在湖和仓之间自由流动。作为新一代大数据技术架构,将逐渐取代单一数据湖和数据仓库架构。湖仓一体兼具数据湖的灵活性与数据仓库的成长性。

-
几款数据湖产品
- Delta Lake
- Iceberg
- Hudi
- Paimon
- 选型对比


-
开放文件格式
- FileFormat:通过 FileFormat 可以将数据格式与 Hive 的每一行 Row 对应起来,形成 Hive 的 Table,这些 Table 的元数据都存储在 Hive 的 MetaData 数据库中。

- TableFormat:Table format 定义了哪些文件构成一张表,这样任何引擎都可以根据 table format 查询和检索数据。Table format 规范了数据和文件的分布方式,任何引擎写入数据都要遵照这个标准,通过 format 定义的标准支持 ACID,模式演进等高阶功能。
- Table format 核心特性
- 结构自由
- 读写自由
- 流批同源
- 引擎平权

-
批流一体和湖仓一体的区别
- OneTable:OneTable 采用源元数据格式并将表元数据提取为通用格式,然后可以将其同步为一种或多种目标格式。思想:OneTable 不是一种新的表格式,而是为 Hudi、Delta、Iceberg 元数据的全向无缝转换提供了所必须的工具和抽象。全向意味着您可以从任一格式转换为其他任一格式,您可以在任何需要的组合中循环或轮流使用它们,性能开销很小,因为从不复制或重新写入数据,只写入少量元数据。在使用 OneTable 时,来自所有 3 个项目的元数据层可以存储在同一目录中,使得相同的 "表" 可以作为原生 Delta、Hudi 或 Iceberg 表进行查询。元数据转换是通过轻量级的抽象层实现的,这些抽象层定义了用于决定表的内存内的通用模型。这个通用模型可以解释和转换包括从模式、分区信息到文件元数据(如列级统计信息、行数和大小)在内的所有信息。除此之外,还有源和目标层的接口,使得其能转入,或从这个模型转出。这些接口允许用户扩展和发展当前 OneTable 为三种主要表格格式提供的功能。例如,开发人员可以实现源层面接口来支持 Apache Paimon,并立即能够将这些表暴露为 Iceberg、Hudi 和 Delta,以获得与数据湖生态系统中现有工具和产品的兼容性。更多详细信息请参考 GitHub 代码库:https://github.com/onetable-io/onetable

- 流批一体是需求,湖仓一体是方案。从 lakehouse 提出的背景看,湖仓一体一定是流批一体,但流批一体不一定基于数据湖,事实上很多传统数仓都具备流批一体的能力。Lakehouse 的设计原则分为功能性设计要素和非功能性设计要素两类。其中,功能性设计要素包括:一体化架构、存算分离、事务和数据一致性、全数据类型。非功能性设计要素包括:弹性高可用、加强的数据治理、尽量少的数据冗余、高并发支持、运维可观测性、高开放性。统一表格式(table format)
湖仓一体的几种方式
- 湖中建仓
- 湖上建仓
- 仓外挂湖
- 比如已经有了 Hive 的数仓存储体系,再引入数据湖的格式,并实现了通过 Hive 对数据湖进行读和写,这种方式就叫做仓外挂湖。
- StarRocks 仓外挂湖是指以 MPP 数据库为基础,使用可插拔架构,通过开放接口对接外部存储实现统一存储,在存储底层共享一份数据,计算、存储完全分离,实现从强管理到兼容开放存储和多引擎。实现方向为增加存储能力,提升查询引擎效率。


常见架构
- Flink + Paimon + OLAP MPP
- MPP + 数据湖存储
集群迁移、升级
迁移、升级工作评估
- 迁移工具
- 迁移方案
- 迁移测试
- 迁移实施
- 一致性
- 稳定性
- 容量、带宽等资源
- 权限
- 迁移过程中的容错保障
- 迁移结果验证
任务迁移、升级
Flink 任务升级
-
- 找到需要停机的任务的 job_id 和 yid
-
- 执行 savepoint 操作,需要带上 job_id,yid 和保存目标路径
-
- 停止旧的任务、提交新的任务时需要指定 -s:savepoint 目录
StarRocks 任务升级(见 StarRocks 官方文档)
DataOps
全方位优化
数据倾斜优化
Hive
- 参数调节:
hive.map.aggr = true在 map 端部分聚合。hive.groupby.mapaggr.checkinterval = 100000在 Map 端进行聚合操作的条目数目。hive.groupby.skewindata=true数据倾斜时负载均衡。
- SQL 语句调节:
- Join 时选择 key 值分布较均匀的表作为驱动表,同时做好列裁剪和分区裁剪,以减少数据量。
- 大小表 join 时,小表先进内存。
- 开启 mapSide join。
- 大表 join 大表时,把 key 值为空的 key 变成一个字符串加上随机数,把倾斜的数据分到不同的 reduce 上,由于 null 值关联不上,因此处理后不影响最终结果。
- 大表 join 大表时,sort merge bucket join。
- 聚合计算依赖的 key 分布不均匀时就会发生数据倾斜,用两次 group by 代替 count distinct 不同指标的 count distinct 放到多段 SQL 中执行,执行后再 UNION 或 JOIN 合并。
set hive.optimize.skewjoin = true;set hive.skewjoin.key = 250000000- Group by 维度过小:采用 sum(),group by 的方式来替换 count(distinct(字段名))完成计算。
Spark
- 使用 Hive ETL 预处理数据
- 过滤少数导致倾斜的 key
- 提高 shuffle 操作的并行度
- 两阶段聚合(局部聚合+全局聚合)
- 将 reduce join 转为 map join(broadcast 大变量)
- 采样倾斜 key 并分拆 join 操作
- 使用随机前缀和扩容 RDD 进行 join
- 自定义分区器
- Spark AQE
Flink
- keyBy 之前发生数据倾斜
- 提高任务并行度
- 使用 shuffle、rebalance 或 rescale 算子将数据均匀分配
- keyBy 后的聚合操作存在数据倾斜
- 使用 LocalKeyBy 的思想:在 keyBy 上游算子数据发送之前,首先在上游算子的本地对数据进行聚合后再发送到下游,使下游接收到的数据量大大减少,从而使得 keyBy 之后的聚合操作不再是任务的瓶颈。(本地聚合攒批之后发往下游)
- keyBy 后的窗口聚合操作存在数据倾斜
- 两阶段聚合
-
- 第一阶段聚合:key 拼接随机数前缀或后缀,进行 keyby、开窗、聚合。注意:聚合完不再是 WindowedStream,要获取 WindowEnd 作为窗口标记作为第二阶段分组依据,避免不同窗口的结果聚合到一起。
-
- 第二阶段聚合:去掉随机数前缀或后缀,按照原来的 key 及 windowEnd 作 keyby、聚合。
-
- 两阶段聚合
背压问题
Flink 背压
-
概述
- 反压场景:系统接收数据的速率高于它处理速率的效率,经常出现在促销、秒杀活动的场景。
- 危害:
- 影响 checkpoint 的时长:checkpoint 时间变长可能导致 checkpoint 超时失败。
- 影响 state 大小:可能拖慢 checkpoint 甚至导致 OOM。
- 反压原理
- TCP 反压:
- TCP 包的格式结构,有 Sequence number 这样一个机制给每个数据包做一个编号,还有 ACK number 这样一个机制来确保 TCP 的数据传输是可靠的,除此之外还有一个很重要的部分就是 Window Size,接收端在回复消息的时候会通过 Window Size 告诉发送端还可以发送多少数据。TCP 就是通过这样一个滑动窗口算法的机制实现 feedback。
-
- 跨 TaskManager,反压如何从 InputChannel 到 ResultSubPartition 中:
- 当 Producer 速率大于 Consumer 速率的时候,一段时间后 InputChannel 的 Buffer 被用尽。InputChannel 都向 LocalBufferPool 申请 Buffer 空间,然后 LocalBufferPool 再向 NetWork BufferPool 申请内存空间。当 Network BufferPool 也用尽的时候,这时 Netty AutoRead 就会被禁掉,Netty 就不会从 Socket 的 Buffer 中读取数据了。过不多久 Socket 的 buffer 也会被用尽(receive buffer),这是 window=0 发送给发送端。这时候 socket 停止发送。
-
- 发送端的 Socket 的 Buffer 也被用尽(send buffer),Netty 检测到 Socket 无法写了之后就会停止向 Socket 写数据。所有的数据就会阻塞在 Netty 的 Buffer 当中,很快 Netty 的 buffer 也不能在写数据了,数据就会积压到 ResultSubPartition 中。和接收端一样 ResultSubPartition 会不断的向 Local BufferPool 和 Network BufferPool 申请内存。Local BufferPool 和 Network BufferPool 都用尽后整个 Operator 就会停止写数据,达到跨 TaskManager 的反压。
-
- TaskManager 内部,反压如何从 ResultSubPartition 到 InputChannel 中:
- 由于 operator 下游的 buffer 耗尽,此时 Record Writer 就会被阻塞,又由于 Record Reader、Operator、Record Writer 都属于同一个线程,所以 Record Reader 也会被阻塞。这时上游数据还在不断写入,不多久 network buffer 就会被用完,然后跟前面类似,经是 netty 和 socket,压力就会向上游传递。
- 缺点:
- 只要 TaskManager 执行的一个 Task 触发反压,该 TaskManager 与上游 TaskManager 的 Socket 就不能再传输数据,从而影响到所有其他正常的 Task,以及 Checkpoint Barrier 的流动,可能造成作业雪崩。
- 反压的传播链路太长,且需要耗尽所有网络缓存之后才能有效触发,延迟比较大。
- 基于 Credit 的反压过程:
- 在每一次 ResultSubPartition 向 InputChannel 发送消息时,都会发送一个 backlog size 告诉下游准备发送多少消息,下游会计算 Buffer 空间大小去接收消息,如果有充足的 Buffer 就返还给上游一个 Credit 告知可以发送消息的大小(图中 ResultSubPartition 和 InputChannel 之间的虚线表示最终还是需要通过 Netty 和 Socket 去通信,并不是直接通信)。
- 优点:
- 基于 credit 的反压过程,效率比之前要高,因为只要下游 InputChannel 空间耗尽,就能通过 credit 让上游 ResultSubPartition 感知到,不需要在通过 netty 和 socket 层来一层一层的传递。
- 另外,它还解决了由于一个 Task 反压导致 TaskManager 和 TaskManager 之间的 Socket 阻塞的问题。
-
定位反压
- 先把 operator chain 禁用,方便定位到具体算子。
- 通过 Flink Web UI 自带的反压监控面板:
- A. 该节点的发送速率跟不上它的产生数据速率。
- B. 下游的节点接受速率较慢,通过反压机制限制了该节点的发送速率。
- 利用 Metrics 定位:
- backpressure Tab 页面
- inpoolUsage=floatingBuffersUsage+exclusiveBuffersUsage
-
反压原因及处理
- 分析方式:
- 使用火焰图分析
- 使用 GC 分析器
- 线程 Dump、CPU profile
- 常见原因及解决方式:
- 数据倾斜
- 增加资源:通过增加资源(例如 TaskManager、CPU、内存等)的方式,来提高整个系统的处理能力,从而降低背压的风险。需要注意的是,增加资源是一种比较暴力的解决方式,并非所有情况都适用。
- 调整拓扑结构:
- 增加缓存队列的长度,以容纳更多的未处理数据。
- 优化算子之间的并行度数量,避免出现单节点瓶颈。
- 上下游并行度保持一致,合并算子链,使用共享资源槽位组。
- 使用窗口(Window)或 State 来帮助管理状态,降低内存占用率。
- 对于大规模任务,可以将任务拆分成多个小任务,以减少单个算子积压数据的风险。
- checkpoint、状态后端等相关优化。
- 外部组件交互:
- 如果发现我们的 Source 端数据读取性能比较低或者 Sink 端写入性能较差,需要检查第三方组件是否遇到瓶颈,还有就是做维表 join 时的性能问题。例如:
-
- Kafka 集群是否需要扩容,Kafka 连接器是否并行度较低。
-
- HBase 的 rowkey 是否遇到热点问题,是否请求处理不过来。
-
- ClickHouse 并发能力较弱,是否达到瓶颈。
-
- 关于第三方组件的性能问题,需要结合具体的组件来分析,最常用的思路:
-
- 异步 io+热缓存来优化读写性能。
-
- 先攒批再读写。
-
- 维表 join 的合理使用与优化。
-
- 如果发现我们的 Source 端数据读取性能比较低或者 Sink 端写入性能较差,需要检查第三方组件是否遇到瓶颈,还有就是做维表 join 时的性能问题。例如:
- 分析方式:
小文件优化
Hive 合并小文件
- 1. 使用命令,自动合并小文件:
- ORC:concatenate;
- Parquet: hadoop -jar parquet-tools-1.9.0.jar merge xxx xxxx
- 2. 调整参数减少 Map 数量:
set hive.input.format=org.apache.hadoop.hive.ql.io.CombineHiveInputFormat;-- 默认set mapred.max.split.size=256000000;-- 256Mset mapred.min.split.size.per.node=100000000;-- 100Mset mapred.min.split.size.per.rack=100000000;-- 100Mset hive.merge.mapfiles = true;set hive.merge.mapredfiles = true;set hive.merge.size.per.task = 256*1000*1000;-- 256Mset hive.merge.smallfiles.avgsize=16000000;-- 16M
- 3. 减少 Reduce 的数量:
- 直接设置 reduce 个数
set mapreduce.job.reduces=10; - 设置每个 reduce 的大小,Hive 会根据数据总大小猜测确定一个 reduce 个数
set hive.exec.reducers.bytes.per.reducer=5120000000;-- 默认是 1G,设置为 5G
- 直接设置 reduce 个数
- 4. distribute by:
insert overwrite table test [partition(hour=...)] select * from test distribute by floor (rand()*5);
- 5. 使用 hadoop 的 archive 将小文件归档:
set hive.archive.enabled=true;set hive.archive.har.parentdir.settable=true;set har.partfile.size=1099511627776;ALTER TABLE A ARCHIVE PARTITION(dt='2020-12-24', hr='12');ALTER TABLE A UNARCHIVE PARTITION(dt='2020-12-24', hr='12');
Spark 合并小文件
- 1. 通过 repartition 或 coalesce 算子控制最后的 DataSet 的分区数(Spark Core)
- 2. Spark AQE
- 3. HINT 方式
- 4. 独立的小文件合并:
set spark.sql.merge.enabled=true;set spark.sql.merge.size.per.task=134217728;
- 5. Spark 自定义异步合并工具类
- 6. 增加 batch 大小(spark streaming)
- 7. 自己调用 foreach 去 append(spark streaming)
Flink 合并小文件
- 1. 自定义 PartitionCommitPolicy:
'sink.partition-commit.policy.kind' = 'metastore,success-file,custom','sink.partition-commit.policy.class' = 'me.lmagics.flinkexp.hiveintegration.util.ParquetFileMergingCommitPolicy'
- 2. Flink 1.12 之后新增的 table/SQL 参数配置:
auto-compaction=truecompaction.file-size=1024M
- 3. Flink StreamingFileSink
- 4. Flink + Hudi:
compaction.max_memorywrite.task.max.sizecompaction.max_memorycompaction.taskscompaction.async.enabled(MOR)compaction.trigger.strategycompaction.delta_commits
- 5. Flink + Iceberg:
.rewriteDataFiles().targetSizeInBytes(128 * 1024 * 1024)
业务层面优化
思考业务逻辑的合理性、可行性、必要性
计算量太大是不是必须的,是否可以减少参与计算的用户量或者时间跨度
计算逻辑是否过于复杂,是否可以简化
根据业务特性做业务数据架构的优化
计算资源优化
开启动态资源分配 dynamicAllocation
CPU/内存配比建议同集群总资源配比,最大化利用集群资源
开启 broadcastjoin,有数大数据平台环境默认关闭(内存限制),spark 官方建议开启,可极大提高 join 小表性能
调节 parallelism 和 repartition 参数,提高并行度
开启 convertMetastoreParquet,充分利用 spark 读 parquet 性能
lateral view explode 优化,多次 explode 前,手动触发 shuffle 操作,减少单分区处理数据量大小
控制输出文件大小,减少小文件数,减轻 nn 压力
充分利用 spark3 AQE 优化
CBO 优化器
Sorted streaming aggregate
Query cache
Spark、Hive、Flink 资源合理分配 shuffle 参数优化
综合优化
存算分离
计算引擎切换
通过 Node label 完成资源隔离
通过 cgroup 完成 CPU 进程的绑定,使其得到充分的利用。
数据结构的合理设计
表设计、存储优化
建表前:表分区、分桶
压缩存储
- 第一级:格式压缩,基于不同表,压缩为 orc/textfile+snappy 格式,压缩率在 30~48% 左右。
- 第二级:使用纠删码技术,从 3 副本减少到 1.5 副本,压缩率在 50% 左右。
建表前:生产级压缩编码设置(ZSTD)
数据写入前:数据布局优化技术
表后期生命周期维护
- 分区 TTL 策略
- 存储冷热分离技术
- 无效表(目录)下线机制
建立数据共享方案
Bitmap
表模型选择
Data Skip
SQL、代码优化
列裁剪,避免 select *
分区裁剪,使用分区字段过滤
条件限制
谓词下推
Map 端预聚合
大 key 的过滤
打散倾斜 key
合适的 join 方式
- Broadcast Hash Join:
- 仅支持等值连接,join key 不需要排序。
- 支持除了全外连接(full outer joins)之外的所有 join 类型。
- Broadcast Hash Join 相比其他的 JOIN 机制而言,效率更高。但是,Broadcast Hash Join 属于网络密集型的操作(数据冗余传输),除此之外,需要在 Driver 端缓存数据,所以当小表的数据量较大时,会出现 OOM 的情况。
- 被广播的小表的数据量要小于
spark.sql.autoBroadcastJoinThreshold值,默认是 10MB(10485760)。 - 被广播表的大小阈值不能超过 8GB。
- 基表不能被 broadcast,比如左连接时,只能将右表进行广播。
- 适合很小的表和大表 join。
- Broadcast 阶段:小表被缓存在 executor 中。
- Hash Join 阶段:在每个 executor 中执行 Hash Join。
- Shuffle Hash Join:
- 选择 Shuffle Hash Join 需要同时满足以下条件:
spark.sql.join.preferSortMergeJoin为 false,即 Shuffle Hash Join 优先于 Sort Merge Join。- 右表或左表是否能够作为 build table。
- 是否能构建本地 HashMap。
- 以右表为例,它的逻辑计划大小要远小于左表大小(默认 3 倍)。
- 适合小表和大表 join。
- 选择 Shuffle Hash Join 需要同时满足以下条件:
- Shuffle Sort Merge Join:
- 仅支持等值连接。
- 支持所有 join 类型。
- Join Keys 是排序的。
- 参数
spark.sql.join.prefersortmergeJoin(默认 true)设定为 true。 - 适合大表和大表的 join。
- Shuffle Phase:两张大表根据 Join key 进行 Shuffle 重分区。
- Sort Phase:每个分区内的数据进行排序。
- Merge Phase:对来自不同表的排序好的分区数据进行 JOIN,通过遍历元素,连接具有相同 Join key 值的行来合并数据集。
- Cartesian Product Join:
- 仅支持内连接。
- 支持等值和不等值连接。
- 开启参数
spark.sql.crossJoin.enabled=true。
- Broadcast Nested Loop Join:
- 支持等值和非等值连接。
- 支持所有的 JOIN 类型,主要优化点如下:
- 当右外连接时要广播左表。
- 当左外连接时要广播右表。
- 当内连接时,要广播左右两张表。
- 总结:
- 优先级为:Broadcast Hash Join > Sort Merge Join > Shuffle Hash Join > Cartesian Join > Broadcast Nested Loop Join.
- 等值连接的情况:
- 有 join 提示(hints)的情况,按照下面的顺序:
- Broadcast Hint:如果 join 类型支持,则选择 broadcast hash join。
- Sort merge hint:如果 join key 是排序的,则选择 sort-merge join。
- Shuffle hash hint:如果 join 类型支持,选择 shuffle hash join。
- Shuffle replicate NL hint:如果是内连接,选择笛卡尔积方式。
- 没有 join 提示(hints)的情况,则逐个对照下面的规则:
- 如果 join 类型支持,并且其中一张表能够被广播(
spark.sql.autoBroadcastJoinThreshold值,默认是 10MB),则选择 broadcast hash join。 - 如果参数
spark.sql.join.preferSortMergeJoin设定为 false,且一张表足够小(可以构建一个 hash map),则选择 shuffle hash join。 - 如果 join keys 是排序的,则选择 sort-merge join。
- 如果是内连接,选择 cartesian join。
- 如果可能会发生 OOM 或者没有可以选择的执行策略,则最终选择 broadcast nested loop join。
- 如果 join 类型支持,并且其中一张表能够被广播(
- 有 join 提示(hints)的情况,按照下面的顺序:
- 非等值连接情况:
- 有 join 提示(hints),按照下面的顺序:
- Broadcast hint:选择 broadcast nested loop join.
- Shuffle replicate NL hint: 如果是内连接,则选择 cartesian product join。
- 没有 join 提示(hints),则逐个对照下面的规则:
- 如果一张表足够小(可以被广播),则选择 broadcast nested loop join。
- 如果是内连接,则选择 cartesian product join。
- 如果可能会发生 OOM 或者没有可以选择的执行策略,则最终选择 broadcast nested loop join。
- 有 join 提示(hints),按照下面的顺序:
用 Distribute By Rand 控制分区中数据量
Group by 优化
Count(distinct)
中间结果的缓存和复用
并行执行
物化视图
Hints 优化
Job 数优化
- 减少 job 数:
- 不论是外关联 outer join 还是内关联 inner join,如果 Join 的 key 相同,不管有多少个表,都会合并为一个 MapReduce 任务。
- JOB 输入输出优化:
- 善用 multi-insert、union all,不同表的 union all 相当于 multiple inputs,同一个表的 union all,相当 map 一次。
Union all 优化
选择合理的排序
SLA 保障
第一,SLA 监控主要监控整体产出指标的质量、时效性和稳定性。第二,链路任务监控主要对任务状态、数据源、处理过程、输出结果以及底层任务的 IO、CPU 网络、信息做监控。第三,服务监控主要包括服务的可用性和延迟。最后是底层的集群监控,包括底层集群的 CPU、IO 和内存网络信息。

时效性的目标也有 3 个,接口延迟的报警、OLAP 引擎报警和接口表 Kafka 延迟报警。拆分到链路层面,又可以从 Flink 任务的输入、处理和输出三个方面进行分析:输入核心关注延迟和乱序情况,防止数据丢弃;处理核心关注数据量和处理数据的性能指标;输出则关注输出的数据量多少,是否触发限流等。


模型优化(设计合理的数仓模型)
是否有现成的数据可以使用或者基于现成的数据进行加工
是否可以将整个计算逻辑进行合理拆分,降低每个子任务的复杂度,
end





