资讯动态

用数据流图搞定ETL:概念、分层与工程实践

发布时间:2026/9/18 13:34:01 来源:尧图企业网站定制
简介围绕ETL流程、数据流图及过程解决方案的PPT课件面向数据仓库开发初学者、ETL工程师以及准备ETL面试的读者。内容完整覆盖ETL定义、实施前提、过程设计原则包括利用数据中转区预处理、主动拉取而非推送、流程化配置管理、数据质量保证等要点系统对比异构与同构两种模式的特点、适用环境与性能差异并针对数据抽取时段冲突、源目标停机、快照定义与重复装载、失败回滚等真实场景给出具体处理思路。资源为1个PPT文件压缩包大小932KB便于下载后快速翻阅适合作为系统学习或面试前查漏补缺的速查材料。已有355人浏览学习对想理解ETL核心流程、模式选型与常见解决方案的读者具有直接参考价值。1. ETL流程和解决方案之间的落差往往差一张数据流图做过数据仓库的人都有同感真正的 ETL流程里写代码的时间占不到一半剩下都在对口径、查数据漂移、跟下游解释为什么这张表又少了 200 行。大部分 ETL 事故追到根上不是 SQL 写错而是设计阶段没把边界、加工职责和存储位置画清楚。数据流图DFD解决的正是这个问题它把「数据从哪来、经过谁、存到哪」用统一符号固化下来评审、排障、交接都靠它。这类主题常以方案 PPT 形态交付但决定方案成色的不是排版而是图和数据流能否经得起追问。下面按「先立概念、再画图、后写代码、最后给交付检查项」展开适合正在设计数仓管道的数据工程师也适合评审别人方案时想快速抓住要害的人。2. ETL流程拆解抽取、转换、加载如何对应数据流图的加工与存储ETL 流程里的每个动作都能在数据流图上找到对应物抽取是加工暂存表是数据存储清洗是另一个加工。先把这层对应关系理顺后面画图、写脚本都只是填充细节。2.1 ETL 概念先立住三个阶段的职责边界ETL 拆开是 Extract、Transform、Load但很多人对三者边界是模糊的。抽取只负责把数据从源系统搬出来不承担业务清洗转换负责标准化、去重、类型统一、代理键生成加载负责把结果写入数仓决定用 insert、update 还是分区交换。边界清楚的最大好处是出问题时能快速定责而不是在一段几千行的 SQL 里翻线索。这三条边界也常成为评审争论点比如「抽取时顺手把空值替换掉算不算越界」。我的口径是抽取阶段只做类型适配凡是依赖业务语义的改动一律留在转换阶段否则边界线形同虚设。增量抽取方式要在概念阶段就定下来时间戳增量要求源表有可靠的增长字段CDC 日志解析会引入 DDL 变更的额外处理触发器方案对源库有侵入。这三种方式决定增量任务的入参怎么设计也直接影响后面能不能断点续跑。全量对比只适合小表数据量上来后性价比下降得很快。选型层面还要处理 ETL 和 ELT 的关系。ELT 把转换推迟到数仓内部执行适合数据量大、计算资源向数仓集中、源库不能承受重查询的场景ETL 则适合敏感数据不允许先落仓、或治理规则要求在下发前完成清洗的场景。不要迷信某一种按「数据量、计算资源、治理要求」三个维度选。2.2 数据流图的分层上下文图、一层图、二层图的分解规则数据流图用四种元素表达系统外部实体、加工、数据存储、数据流。外部实体是数据来源和去向加工是处理动作数据存储是静止的数据数据流是带名字的箭头。ETL 场景里源库和报表平台是外部实体清洗逻辑是加工ODS 表、DWD 表是数据存储搬运动作就是数据流。数据流图的分解是画图的核心手法本质上是上下文数据流图的逐层展开。第 0 层上下文图只有一个加工代表整个 ETL 系统外部实体围绕在四周第 1 层把这个加工分解成抽取、转换、加载等主加工第 2 层再把复杂加工继续展开。分解有三条硬规则一是平衡父图加工的输入输出数据流必须与子图边界流一致二是每条数据流都要命名不能出现无名箭头三是加工不能只有输入没有输出反之亦然。分解的停止条件比规则本身更容易被忽略。我的判断标准是子加工能不能对应到一个能独立编码、独立测试的函数或 SQL 块。能就停止不能就继续拆。拆到「一个人一天能写完并自测」的粒度图才不会变成一张谁都看不懂的蜘蛛网。2.3 ETL 步骤到数据流图元素的映射附平衡检查表ETL 环节数据流图元素画法要点源业务库外部实体只标数据流向不画内部结构增量抽取任务加工输入源表数据流输出到暂存存储暂存区ODS/Staging数据存储必须有流入和流出两条流清洗转换任务加工第 2 层展开成字段级子加工数仓明细表数据存储只允许加载加工写入下游报表/应用外部实体消费方不需要回写路径画完图按这张表自检三个点每个加工是否至少有一条输入和一条输出每个数据存储是否同时被读和写所有数据流的名字能否在表或字段层面找到对应物。很多评审只盯着业务逻辑对不对反而漏掉这些结构性约束。结构错误不报错只会让图永远停在「看起来对、实际上没人能照着实现」的状态。结构检查可以用脚本固化。把每个加工输入输出的数量登记到元数据表一条 SQL 就能找出孤立加工-- 在元数据表 dfd_process_meta 上做结构自检 -- in_flow_count / out_flow_count 由画图工具或评审录入生成 SELECT process_id, process_name FROM dfd_process_meta WHERE in_flow_count 0 OR out_flow_count 0 OR in_flow_count IS NULL OR out_flow_count IS NULL;这条查询的作用是把「加工必须有进有出」变成持续执行的约束。把它放进评审的前置检查脚本里谁提交的图上出现孤立加工当场打回比人肉看图高效得多。3. 画数据流图记号选择与 ETL 场景的分层画法图不是画得越细越好而是画到「能照着实现、能经得起追问」为止。下面按记号选择、上下文图、分层展开、脚本校验四步走每步都给出可以照做的判断标准。3.1 先选记号Yourdon 与 Gane-Sarson 的取舍DFD 有两套流传最广的记号。Yourdon 记号里加工画成圆数据存储画双竖线Gane-Sarson 记号里加工画圆角矩形数据存储画开口矩形。两套表达力等价但混用会让评审反复解释。常见的画图工具对两套都支持拖拽即可关键是全图统一。元素Yourdon 记号Gane-Sarson 记号外部实体矩形矩形加工圆形圆角矩形数据存储双竖线开口矩形数据流带箭头直线带箭头直线我习惯用 Gane-Sarson因为圆角矩形里能放下「加工编号 动词短语」两行字可读性更好。如果团队已有存量图以既有记号为基准统一比新颖重要。3.2 第 0 层上下文图把 ETL 系统的边界划在源和仓之间上下文图的目标是定边界。图里只有一个加工名字就是整个 ETL 系统。围绕它的外部实体通常有四类源业务系统、数仓存储集群、下游数据应用、运维触发方。网上能搜到不少教务管理系统数据流图、虚拟桌面 VDI 数据流图之类的练习素材但 ETL 上下文图的难点不在外框画法而在两个容易被忽略的判断。第一个判断是「运维调度要不要画进来」。DFD 只表达数据不表达控制调度信号属于控制流硬画成数据流会破坏图的一致性。常见做法是把调度画成外部实体连向系统的触发流注明是控制事件调度本身的时序放到单独的任务设计文档里表达。第二个判断是数据流上标注什么粒度。上下文图只标表名或数据量级即可字段级细节放到第 2 层概览和明细混在一张图上谁都看不下去。3.3 第 1 层图与第 2 层图把转换加工展开到能直接编码第 1 层图把唯一的加工展开成三个主加工抽取、清洗转换、加载并画出暂存区和明细表两个数据存储。到了第 2 层把「清洗转换」继续拆成字段标准化、去重、代理键生成、SCD 处理等子加工。注意每级展开都必须满足父图子图平衡否则第 2 层就是无效的。例如「订单金额标准化」子加工输入是暂存区原始金额流输出是清洗后的金额流加工内部要处理单位不统一、负金额、精度溢出三种情况。展开到这个粒度写代码的人就不需要回头猜图了。3.3.1 查询修改数据流图时先查这三处评审和修改已有数据流图问题最多出现在三处。一是平衡被破坏改了父图加工的数据流子图没同步。二是命名偷懒出现「数据」「结果」这类万能名无法对应具体字段。三是存储读写方向缺失数据存储只画了写入没画读取或反过来后者直接暴露「只写不读」的无效存储。把这三条写成评审检查表画图就从个人风格变成团队工程规范。3.4 用 Python 校验父图与子图的平衡手动检查平衡在高层图还能应付加工编号深到 1.3.2 之后必须靠脚本。下面这段代码用最简模型实现每个加工记录输入输出流子图内部流成对出现自动抵消剩余外部流必须与父图完全一致。# dfd_balance.py —— 校验父图与子图的数据流平衡 # 数据流名需全局统一内部流会同时出现在某子加工输出与另一子加工输入 processes { 1: {in: {源表数据}, out: {仓内表数据}}, 1.1: {in: {源表数据}, out: {暂存区数据}}, 1.2: {in: {暂存区数据}, out: {清洗后数据}}, 1.3: {in: {清洗后数据}, out: {仓内表数据}}, } def outer_flows(children: list[str]) - tuple[set, set]: ins, outs set(), set() for cid in children: ins | processes[cid][in] outs | processes[cid][out] internal ins outs # 成对出现的内部流抵消 return ins - internal, outs - internal def check(parent_id: str, children: list[str]) - None: pin, pout processes[parent_id][in], processes[parent_id][out] cin, cout outer_flows(children) assert pin cin, f{parent_id} 输入不平衡: {pin ^ cin} assert pout cout, f{parent_id} 输出不平衡: {pout ^ cout} print(f{parent_id} 平衡检查通过) check(1, [1.1, 1.2, 1.3])代码逻辑子图里「暂存区数据」「清洗后数据」同时是一个子加工的输出和另一个子加工的输入属于内部流取交集后从两侧集合剔除剩下的就是子图对外的边界流必须与父图完全一致。实际项目中可以把这张表换成元数据库里的两张表用 SQL 集合运算做同样的差集检查跑一遍比十个人轮流看图可靠得多。4. ETL过程解决方案工具选型、调度设计与可重跑的 Spark 脚本解决方案落到代码之前有两件事必须先定用什么工具承载调度与计算失败之后怎么重跑才算安全。这两件事定不下来脚本写得再好也守不住交付进度反过来选型对了、重跑策略理顺了实现本身反而简单。4.1 常见 ETL 工具分三类先分清要解决哪个问题「ETL 工具」其实覆盖三类不同的东西选错类型是方案返工的第一大原因。类型代表擅长不擅长调度编排型Airflow、DolphinScheduler任务依赖、告警、失败重跑复杂转换逻辑数据整合型Kettle、DataStage、云上数据集成服务可视化拖拽、异构数据源连通海量数据吞吐、细粒度调优计算引擎型Spark、Flink大数据量转换spark etl 脚本是高频实践自身不解决调度需配合编排框架我的选型建议是调度编排与计算引擎分离。团队如果是 SQL 熟练工「调度框架 SQL/Spark 脚本」性价比最高如果业务人员要频繁自助接数数据整合型工具能降低接入门槛但要接受它在复杂字段映射和血缘追溯面前的吃力。4.2 最小可运行的 Spark ETL 脚本与参数说明典型场景每天从 MySQL 业务库抽取订单增量清洗后按业务日期分区写入数仓明细表结束时更新水位表。下面是可运行骨架# spark_etl_daily.py —— 订单日增量 ETL from pyspark.sql import SparkSession from pyspark.sql.functions import col, to_date, when, lit, max as _max from pyspark.sql.types import DecimalType spark SparkSession.builder \ .appName(order_daily_etl) \ .config(spark.sql.shuffle.partitions, 200) \ .config(spark.sql.sources.partitionOverwriteMode, dynamic) \ .enableHiveSupport() \ .getOrCreate() # 抽取从上批次水位开始读增量水位表就是一个数据存储 wm spark.read.parquet(/ods/order_watermark).first()[0] src spark.read.format(jdbc) \ .option(url, jdbc:mysql://10.0.1.5:3306/biz) \ .option(dbtable, f(SELECT * FROM t_order WHERE update_time {wm}) t) \ .option(driver, com.mysql.cj.jdbc.Driver) \ .load() # 转换主键去重、金额空值兜底、统一小数精度、生成分区字段 clean src.dropDuplicates([order_id]) \ .withColumn(amount, when(col(amount).isNull(), lit(0)) .otherwise(col(amount)).cast(DecimalType(12, 2))) \ .withColumn(biz_date, to_date(col(create_time))) # 加载动态分区覆盖只替换本次涉及的业务日期分区 clean.write.mode(overwrite) \ .format(parquet) \ .partitionBy(biz_date) \ .save(/dwd/order_detail) # 更新水位供下一个增量批次取数 clean.select(_max(update_time).alias(wm)) \ .write.mode(overwrite).parquet(/ods/order_watermark)参数说明盯三个点。第一partitionOverwriteModedynamic必须显式打开否则mode(overwrite)会清空/dwd/order_detail的整个表目录这是增量任务里最典型的误操作还原成本极高。第二shuffle.partitions设 200 是起步值建议按「目标分区数 × 2 到 3」调整设太大会写出大量小文件设太小则单 task 内存吃紧。第三水位写在加载之后且水位表本身就是数据流图里的一个数据存储脚本体现的正是「源 → 加工 → 存储 → 加工 → 存储」的完整闭环。脚本跑失败时的排查顺序也有讲究先看水位对不对再看源表字段有没有变更最后才翻 task 日志。很多所谓的神秘失败其实是源表加了一列之后 JDBC 映射错位。4.3 调度与异常重跑数据流图之外必答的环节调度依赖是 ETL 成败的隐形层。层与层之间要声明「上游分区就绪」依赖不是只看任务是否成功。常见做法是上游成功后校验目标分区文件数或行数通过才允许下游启动。空表不告警直接放行会让下游在静默中被写空这类事故占日常排障的大头。失败场景处理方式抽取失败直接重跑读取逻辑幂等不用清理暂存转换失败先看脏数据样例修规则后整层重算重跑前清本次半成品加载失败先查目标分区是否残留半截数据按分区清理后再写上游数据回溯从受影响的最上游任务级联重跑禁止只刷下游重跑的核心是幂等同一批次跑两次结果必须与跑一次完全一致。「先清目标分区、再写、最后更新水位」这个顺序不能乱乱了就会出现「批次 2 已经写完批次 1 又覆盖回来」的现场。告警还要分级空表、行数骤减走电话偶发重试成功走群通知。全部问题都打同一个渠道值班的人很快就麻了。4.4 两道常考的 ETL 面试题顺便把边界讲清数据流图和 ETL 是两类高频面试题下面两道能检验上面的内容是否真懂。第一道ETL 和 ELT 的选型依据是什么。回答层次应落在「转换发生在源端还是仓内取决于数据量、源库负载承受力、敏感数据是否允许先落仓」能对着具体场景给理由比背定义有用。第二道增量抽取有哪几种方式各自的问题是什么。时间戳要求源表有可靠的增长字段CDC 日志解析要处理 DDL 变更触发器和全量对比适合小表。回答时补一句「水位管理和断点续跑才是增量方案的核心难点」面试官基本就能判断你踩过坑。5. 交付前的最后一道检查行数对账、水位减一秒与断点续跑演练上线前把检查项做成任务的后置动作比靠人肉盯日志可靠得多。检查项不需要多行数对账、水位推进、断点续跑这三件事守住绝大多数增量 ETL 的夜间事故都能被挡在告警之前。5.1 行数对账加载完成后的第一个后置动作源表当日增量与数仓目标表当日分区做行数对比。对账 SQL 放在目标侧执行源侧走只读副本避免给生产库加压力-- post_etl_check.sql —— 加载完成后由调度后置节点执行 SELECT (SELECT COUNT(*) FROM biz_replica.t_order WHERE update_time 2025-01-01) AS src_cnt, (SELECT COUNT(*) FROM dwd.order_detail WHERE biz_date 2025-01-01) AS dwd_cnt;比对结果允许src_cnt dwd_cnt的小幅偏差偏差来自抽取窗口内源库仍在写入走副本的场景偏差可忽略。如果dwd_cnt明显偏小优先怀疑分区覆盖模式配错或水位推进过快而不是源库数据问题。5.2 水位减一秒增量 ETL 里最容易被低估的细节水位取max(update_time)有一个隐蔽问题时间戳精度有限时同一秒内前一批次已读到部分数据后一批次从这一秒继续读就会漏掉同秒后写入的行。常见做法是水位向前回拨一个时间精度即取max(update_time) - interval 1 second作为下一批次的起点再由脚本里的dropDuplicates([order_id])把重复行吸收掉。这个技巧成本一行能免掉大部分「少了几行又找不到原因」的深夜电话。断点续跑验证建议上线前做一次破坏性演练在加载阶段中途强制终止任务检查目标分区是否残留半截数据修复后重跑再对比演练前后的行数和金额汇总。能通过这次演练重跑按钮才敢交给值班的人。把上面的检查清单固化成post_etl_check.sql连同平衡校验脚本挂到调度任务的后置节点然后跑一遍完整演练这套 ETL 过程方案才算真正交付。本文还有配套的精品资源点击获取

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价