免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

数据流图驱动的ETL流程设计:从数据抽取到加载的工程实践

数据流图驱动的ETL流程设计:从数据抽取到加载的工程实践 简介围绕ETL全流程的PPT课件面向数据仓库开发人员、运维工程师及正在准备ETL面试的求职者重点解决抽取、转换、加载过程中的方案选型与异常处理问题。内容系统梳理了ETL定义、实施前提范围确定与工具选型、核心原则数据中转区预处理、主动拉取、流程化配置、数据质量保证并对同构与异构两种模式从性能、环境、维护、灵活性和排错难度等维度作出详细比较同时结合抽取时间区间、源库生产时段、系统Down机、快照机制、主外键关联等实际场景分析注意事项给出抽取失败后重新抽取、装载时回滚、按主键判重更新、文本文件快照核查等错误处理思路也覆盖了抽取分析、变换数据、装载数据、数据质量控制四大解决方案步骤。资源为1个PPT演示文稿压缩包约932KB体量精简但知识密度高ETL数据流图、架构模式对比和典型排错流程均在一个课件内集中呈现适合面试前快速回顾或作为项目方案参考。已有355人学习下载。1. ETL 流程不是“先做数仓”而是“先画清楚数据怎么走的”做数据接入的这些年我见过最多的返工不是因为 ETL 逻辑写错而是因为还没弄清楚源头数据从哪来、经过哪些加工、最后落到哪张表就开始写脚本。一个订单系统每天产生几千万行明细销售部门只要汇总数财务却要核对到单据级别这两条数据流如果不提前分开后面无论是用 Spark 还是用存储过程都会在口径上反复扯皮。ETL 流程解决的是“怎么从源头把数据加工成可用的状态”数据流图解决的是“这个可用的状态到底长什么样、中间经过哪几个加工节点”。把数据流图当作 ETL 的第一步不是画给管理层看的而是给每一个转换脚本定义输入输出边界。我一般会先画上下文数据流图再逐层分解直到每个加工节点都能对应到一个 ETL 任务。这个习惯能让团队里刚接手的新人也敢改代码因为他能从图上看出某个字段是从哪张表带入的改坏了会影响哪个下游。如果你是要准备 ETL 面试题这也是最容易被考到的一层——面试官往往不关心你背了多少 ETL 概念而会随手画一个数据流场景看你能否把它拆成可实现的抽取、转换、加载步骤。2. 数据流图要画出什么上下文图、分解图与 ETL 边界2.1 上下文数据流图把系统当黑盒先定外部实体上下文数据流图是数据流图的第 0 层它只画一个系统框以及和系统交互的外部实体。外部实体是数据产生方或消费方比如教务管理系统里的学生、教师、教务员或者更实际一点源数据库、上游文件服务器、下游数据仓库。不要把表名画在这一层更不要把 ETL 任务里的临时表写上去否则就失去了上下文图的概括作用。还有人把用户操作流程和系统批处理流程混在一张上下文图里结果图里出现一堆人工节点这与 ETL 场景并不对应反而干扰了边界定义。对于教务管理系统数据流图这类典型业务系统外部实体通常是学生、教师、教务管理员系统内部是一个黑盒。我要做数仓时需要把“教务系统”框拆开看成选课流水、成绩表、教室资源等几个逻辑实体。这一步的关键是确定数据流的方向是外部流向系统还是系统流向外部录入接口和报表查询的方向完全不同。上下文图定好之后才能聊第二个问题——数据要不要进数仓以及进数仓之前是否需要加工。2.2 上下文数据流图的分解按业务事件拆子过程上下文图只有一层不能满足 ETL 设计。需要做上下文数据流图的分解把系统框按业务事件逐步展开先画 1 层图列出主要处理过程例如“1.1 校验选课数据”“1.2 计算成绩汇总”“1.3 生成班级报表”再对每个过程往下拆直到过程内部没有复杂逻辑只剩简单的输入输出。每个过程都有自己的编号数据流也要带名字和属性列表。养成这个习惯后当别人拿着一张旧图要查询修改数据流图时你可以通过过程编号快速定位到对应的抽取规则。这里需要建一张映射表把数据流图的层次对应到 ETL 工程产物这样团队协作时不会各画各的。数据流图层次典型内容对应 ETL 产物第 0 层上下文图外部实体、系统框数据源清单、目标表清单第 1 层处理器图主要业务过程、数据流映射文档、转换逻辑清单第 2 层详细图单个过程细节、数据字典抽取 SQL、转换脚本在实际项目中数据流图上的每个过程框我都会绑定一个过程编号这个编号会出现在 ETL 任务名和调度日志里。比如“1.2 计算成绩汇总”对应etl_process_1_2排错时看日志就能知道是哪个业务环节出了问题不需要去翻代码里的业务注释。2.3 从数据流图推导 ETL 抽取与加载的边界数据流图和第 2 章讲的数据字典表核心价值是把“图”变成“可查询的元数据”。我一般在数据流图完成后会同步维护一张dataflow_dict表记录每条数据流的编号、字段组成、经过的处理过程编号。这样当需求变更是查询修改数据流图就不需要再打开设计稿直接用 SQL 就能定位影响面。一个简化版本的数据字典查询如下SELECT d.flow_id, d.process_id, d.field_name, p.process_name FROM dataflow_dict d LEFT JOIN process_dict p ON d.process_id p.process_id WHERE d.flow_id FLOW_STU_SCORE ORDER BY d.process_id, d.sort_no;这里的flow_id对应数据流图中一条箭头的名字process_id对应处理过程编号。执行这条 SQL 后你可以看到成绩数据流从源到目标经过了哪些加工节点、涉及哪些字段。只要这个字典维护得干净后面做 ETL 任务拆分时几乎不用再讨论“这个字段要不要清”这种问题因为数据流图已经把边界画死了。3. 从数据流图到 ETL 流程设计抽取、转换、加载的参数与实现3.1 抽取全量、增量与 CDC先定水位线进入 ETL 流程设计时先要处理抽取策略。数据流图上的每个数据流都要回答三个问题源端是否支持只读业务的变更时间字段是什么允许的延迟是多长一般支持时间戳的源表优先做增量抽取用updated_at作为水位线没有时间戳但有主键的表要么全量抽取要么做基于日志的 CDC。水位线的位置很关键太靠前漏数据太靠后重复数据。我通常会把水位线从源表单独抽出来存一张控制表而不是在抽取 SQL 里硬编码时间。一个常见的增量抽取 SQL 如下-- 用控制表里记录的水位线过滤增量数据 SELECT order_id, user_id, amount, status, updated_at FROM source_order WHERE updated_at (SELECT last_watermark FROM etl_watermark WHERE table_namesource_order) AND updated_at CURRENT_TIMESTAMP;last_watermark是上一次成功任务写入的水位线CURRENT_TIMESTAMP是本次调度的时间。这里使用左开右闭区间保证和上次边界不重叠也不漏数据。抽取完成后要在同一个事务里更新水位线否则任务重跑会产生重复数据。这个细节很少有人提但它正是 ETL 面试题里“如何保证增量任务不重不漏”的标准答案。对于维度表这类数据量小、更新频率低的表我一般选择全量覆盖省去维护水位线的成本。抽取频率也不是越短越好比如一个 5 分钟产生一次的日志表如果目标只是天级报表每小时抽一次就够了。数据流图上目标存储的“刷新频率”其实已经写清楚了照着定调度即可。3.2 转换清洗、映射、规范化用 Spark ETL 脚本实现转换是 ETL 流程里最容易被业务复杂拖垮的部分。根据数据流图上的每个处理过程通常要完成四类动作清洗空值、非法字符、映射编码键到维度键、规范化统一时间格式、货币单位、聚合按维度汇总。对于中大规模数据我一般用 Spark ETL 脚本做批处理因为它能把清洗和聚合放在同一份代码里且容易重跑。一个简化的 Spark ETL 脚本骨架from pyspark.sql import SparkSession, functions as F spark SparkSession.builder \ .appName(etl_dim_student) \ .config(spark.sql.shuffle.partitions, 20) \ .getOrCreate() # 从源库读取学生信息表url 中的参数从配置中心注入 df spark.read.format(jdbc).options( urljdbc:mysql://source_host:3306/campus, dbtablestudent_info, user${MYSQL_USER}, password${MYSQL_PASSWORD} ).load() # 清洗去掉软删除和主键为空的数据 df df.filter(F.col(status) ! deleted) \ .dropna(subset[student_no]) \ .withColumn(gender_code, F.when(F.col(gender) 男, M).otherwise(F)) \ .withColumn(etl_time, F.current_timestamp()) # 覆盖写全量维表 df.write.mode(overwrite).format(parquet).saveAsTable(dw.dim_student)参数说明spark.sql.shuffle.partitions控制聚合和关联时 shuffle 分区的数量20 只适合小维表真正跑大流量需要按数据量调到 200 到 1000。dropna(subset[student_no])用于去掉主键空值避免加载后出现无法关联的孤儿行。gender_code把业务值映射成目标端标准编码etl_time则是数据流图中“加载时间”的落地标记。写模式使用overwrite因为维表全量重刷比增量合并更容易排查。这些配置全部外部化到环境变量改动业务逻辑时不需要重新提交 Spark 任务。简单映射和过滤用 SQL 更好当场就能对着数据流图核对结果需要跨系统连接、窗口函数或复杂标准化的场景才用 Spark 脚本因为测试成本低、调试也方便。3.3 加载覆盖写、追加与 SCD控制幂等加载策略需要和数据流图的存储层约定一致。数据流图在画目标数据存储时最好就标上是“每日快照”还是“累计流水”这直接决定加载方式。覆盖写适合维度表或每日全量事实表追加写适合日志型流水SCD 缓慢变化维则要对代理键做 update/insert。下表是我常用的选型参考场景写入模式幂等性每日全量维度表overwrite分区高删完重写增量事实表append 到分区中依赖上游去重SCD 维度表merge / upsert高靠自然键去重加载时要特别注意分区边界。我用 Spark 写数据时会把数据流图中的“业务日期”字段作为分区键写入而不是按执行日期分区。这样重跑昨天的数据只要覆盖对应分区不会污染今天的结果。对于 append 模式最好在写入前先对目标分区做去重或者用一个临时表先去重再插入否则任务失败重跑会出现重复行。针对 Hive 数仓加载 SQL 要写成分区级覆盖INSERT OVERWRITE TABLE dw.dws_student_score_di PARTITION (dt 2024-06-01) SELECT student_no, course_id, score, score_type FROM staging.score_cleaned WHERE dt 2024-06-01;注意这里的OVERWRITE只覆盖dt2024-06-01不会影响其他分区。如果不带分区字段直接INSERT OVERWRITE会清空整张表这是很多线上事故的来源。加载完成后还要更新控制表里的水位线和执行状态把这一步放在同一个调度任务里才能保证整个 ETL 流程可重跑、可追踪。4. 一套可落地的 ETL 过程解决方案调度、监控与重跑4.1 选型从 ETL 工具到调度框架当数据流图和 ETL 任务拆分都清晰后就是选型和落地。ETL 工具目前选择很多Kettle、DataX 适合小团队和数据库间同步dbt 适合 SQL 优先的转换层Spark 适合 PB 级批处理。调度层面我一般用 Apache Airflow 或 DolphinScheduler因为任务依赖和失败告警天然是它们的核心功能。但选型的原则不是哪个框架代码更酷而是看它是否满足三点任务能否表达依赖关系、失败能否被观察到、重跑是否方便。如果你只有十张表用带锁的 Shell 脚本加 crontab 也能撑住有几十张表且依赖链复杂才值得上调度平台。我见过很多团队一上来就搭 Airflow结果 py 文件里只有一堆没有任何依赖的BashOperator调度器变成了定时器反而是负担。评分标准应该是先用最小的成本把 ETL 流程跑起来再按需引入框架。4.2 解决方案的最小骨架分区表 任务编排 失败重试一个最小但完整的 ETL 工程至少要包含一个带锁的 Shell 调度脚本。下面这个脚本用文件锁防止任务重入并按日期参数跑完“抽取、转换、加载检查”三个阶段#!/usr/bin/env bash set -euo pipefail source /etc/etl_profile.sh # 日志目录按日期建好 log_dir/var/log/etl/$(date %F) mkdir -p $log_dir # 用 flock 防止上一轮还没跑完就重入 exec 9$log_dir/order_etl.lock if ! flock -n 9; then echo previous job still running, exit $log_dir/run.log exit 1 fi # step 1: 抽取 python3 ${ETL_HOME}/extract_order.py --date $1 # step 2: 转换 spark-submit --master yarn --deploy-mode cluster ${ETL_HOME}/transform_order.py --date $1 # step 3: 加载并做门禁检查 mysql -h ${DW_HOST} -u ${DW_USER} -p${DW_PASS} dw ${ETL_HOME}/load_check.sql echo $(date %F_%T) order etl done $log_dir/run.logset -euo pipefail保证任一步出错立即退出flock锁解决调度重入问题比如上一个任务还没跑完下一个调度周期已经到了此时直接退出而不是并发跑。这里最关键的一点是$1作为业务日期参数贯穿三个步骤所有脚本只认同一个日期避免抽取、转换、加载各自取当前时间导致跨天不一致。这些脚本挂到 Airflow 的BashOperator上也能直接用不需要改写业务逻辑。失败重试要有上下限。我一般把依赖上游源的抽取任务重试 3 次每次间隔 5 分钟转换和加载任务不自动重试因为如果是数据逻辑问题重跑多少次都一样。自动重试只应针对网络瞬时故障和资源竞争业务逻辑错误必须留给人来处理否则日志里全是同样的错误堆栈。4.3 数据质量检查与血缘追踪ETL 过程不能“跑完就算成功”还要在加载前做数据质量门禁。常见做法是一组检查 SQL在数据写入目标表之前先写入临时表用失败条件使任务退出。例如检查目标表主键是否重复SELECT COUNT(*) AS dup_cnt FROM ( SELECT dw_id, COUNT(*) AS c FROM staging.score_cleaned GROUP BY dw_id HAVING COUNT(*) 1 ) t;如果dup_cnt 0脚本应该停止加载。配合数据流图还可以把检查结果写入一张血缘表记录某个表的数据来自哪个process_id。排错时拿到一条异常记录能直接从目标表反查到源表。下面这几类检查项是正式方案里必备的检查项SQL 判断条件失败动作主键重复dup_cnt 0终止任务并锁定分区空值比例空值行 / 总行数 0.05邮件告警允许继续金额负值MIN(amount) 0终止任务分区延迟最新分区时间早于调度时间告警不阻塞即使没有专门的数据质量工具用这些脚本也能形成“任务级门禁”让 ETL 流程在数据错误进入目标表之前就被拦下来。5. 用数据流图做 ETL 覆盖矩阵一条条核对变更影响5.1 覆盖矩阵怎么建在项目交接和维护阶段我最常做的一件事是把数据流图翻译成覆盖矩阵。做法是从数据流图里抽出所有数据流和处理过程逐行登记到表格然后在每一行后面标注它对应的 ETL 任务编号、调度周期、检查 SQL。这个表本质上就是把图上的线与实际代码连接起来作为变更影响面分析的单据。数据流编号来源处理过程目标ETL 任务调度检查点FLOW_STU_SCORE教务库.score1.2 成绩标准化dw.dws_scoreetl_score_dailydaily 03:00重复主键、空值FLOW_STU_INFO教务库.student1.1 维度清洗dw.dim_studentetl_student_hourlyhourly地址字段空置率每一行都对应数据流图中的一个箭头和一个过程框。建好矩阵后如果需求方说“成绩表要加一个学分字段”我可以先查FLOW_STU_SCORE涉及的处理过程和 ETL 任务再从任务代码里定位到转换脚本的哪一行而不是从数仓底层一张一张表翻。5.2 需求变更时怎么查询修改数据流图并同步矩阵数据流图不是一次性产物业务变化后必须同步修改。正确的修改路径是先在数据流图中找到受影响的子过程和数据流用数据字典查询它绑定的字段和下游节点然后修改数据字典和覆盖矩阵最后才动 ETL 脚本。具体我一般按下面几步走用数据流编号在dataflow_dict表查出现有的字段清单。在图上把新增或删除字段的箭头画出来注意不能并列画多条不连通的数据流。更新process_dict中的过程逻辑描述并修改覆盖矩阵中的目标表和检查点。改 ETL 脚本重跑受影响分区用前面的检查 SQL 验证行数和字段值。这四步走完数据流图、覆盖矩阵和线上代码才是一致的。如果需求变更只改代码不改图下次做影响分析就会漏掉这个节点。5.3 面试问答里的运用方式如果你正在准备 ETL 面试题不要只背概念。面试官往往会给你一张简化的 DFD问你会怎么实现。你可以直接用这套覆盖矩阵来回答先确认数据流图的边界再说出哪些数据流对应抽取、哪些过程对应转换、哪些目标存储对应加载最后指出哪一步失败会影响下游。这样的回答既能体现对数据流图的理解又能落到可执行的 ETL 方案上比空谈“先抽取再转换”更有说服力。维护好这张覆盖矩阵后续每一次需求变更都有据可查线上问题也能一路回溯到源头。本文还有配套的精品资源点击获取
返回列表