前段时间我拿一个银行场景练了练数据抽取,任务本身不复杂:从几个业务库里把客户表和交易流水表抽到数仓的 ODS 层,再用 DolphinScheduler 做成每天自动跑的调度。真正动手之后才发现,数据抽取的难点根本不在 SQL 怎么写,而在字段映射、增量边界、调度依赖和排错链路这些不起眼的细节上。这篇文章把整个练习的完整过程拆开讲,适合刚接触数据抽取、想用 DolphinScheduler 做数据库同步的读者,也适合准备用银行类数据做项目练习的人参考。我会把当时为什么这么设计、配置时踩了哪些坑、最后怎么稳定跑起来,都摊开来说清楚。
1. 银行数据抽取练习到底在练什么
1.1 为什么选“银行场景”做数据抽取练习
银行类数据是数据抽取里很典型的一类练习对象。一方面它的表结构有一定的业务含义,客户表、账户表、交易流水表、渠道日志表,字段多、类型杂,比单纯用 t_user 和 t_order 练手更接近真实情况;另一方面银行数据对完整性和时效性要求高,日终跑批、增量同步、断点续跑这些概念都能在练习里自然带出来。
做这个练习之前,我建议你先别急着装工具、写 SQL,而是想清楚一个问题:这个练习到底想练什么?如果只是想学会“把表 A 的数据复制到表 B”,那用 Navicat 或者一条 INSERT INTO ... SELECT 就够了。但“数据抽取”真正的考点在于:源表数据会变,目标表要怎么跟着变?同一个字段在源库和目标库类型不一致怎么办?今天跑了明天的数据会不会重复?调度挂了之后怎么恢复?
这些问题恰恰是银行小项目练习里最容易遇到的。我给自己定的目标也很简单:用 DolphinScheduler 把两个业务库里的 4 张表,按天增量抽取到一个统一的 ODS 库,支持失败重跑,跑完后能看到明确的成功和失败记录。能把这个闭环跑通,比盲目追求复杂架构有用得多。
1.2 一个能跑通的最小闭环:源表、目标表和边界
练习范围一定要控制住,不然很容易陷进去。我第一次做的时候想一口气抽 10 张表,结果光梳理表关系就花了半天,最后真正调通的只有 2 张。后来我重新做了边界设计:源库用 MySQL,目标库也用 MySQL,这样连接和类型转换最简单;只抽两张有代表性的表,一张是变化不频繁的客户维表,一张是每天都会新增的交易流水表。
源表我模拟的是银行核心系统导出后的业务库,表结构大概长这样:
-- 客户表 CREATE TABLE customer ( cust_id VARCHAR(32) PRIMARY KEY, cust_name VARCHAR(128), id_type VARCHAR(16), id_no VARCHAR(64), mobile VARCHAR(20), create_time DATETIME, update_time DATETIME ); -- 交易流水表 CREATE TABLE txn_log ( txn_id BIGINT PRIMARY KEY AUTO_INCREMENT, cust_id VARCHAR(32), txn_type VARCHAR(8), txn_amount DECIMAL(18,2), txn_time DATETIME, create_time DATETIME, KEY idx_txn_time (txn_time) );目标库就是数仓的 ODS 层,我建了一个独立 schema,表名加了ods_前缀,字段比源表多两个:etl_time记录抽取时间,data_date记录业务日期。这样每次跑批,数据属于哪一天一目了然,也方便后面做数据回溯。
最小闭环的边界就是:每天凌晨 1 点,DolphinScheduler 触发工作流,把前一天的全量客户数据(用 update_time 判断变化)和前一天新增的交易流水(用 txn_time 判断)抽到 ODS 库,最后发一条任务成功或失败的通知。不涉及分库分表、不涉及 Kafka、不涉及宽表加工,先把主链路跑通再说。
1.3 技术选型:为什么用 DolphinScheduler 直接抽数据库
市面上的数据同步工具很多,DataX、SeaTunnel、Flink CDC、Canel 各有各的适用场景。但在练习场景里,我选了 DolphinScheduler 做调度,配合它自带的 SQL 节点直接完成“查源库、写目标库”的动作,没有额外引入同步引擎。为什么?
原因是这个练习的核心目标是理解调度编排和数据抽取的关系。DolphinScheduler 的 SQL 节点支持选择数据源、执行查询、将查询结果直接落到目标库,天然适合“数据库到数据库”的抽取。相比再套一个 DataX,虽然性能上限更高,但配置链路长了一截,练习时容易把注意力全放在工具配置上,反而忽略了增量条件、幂等设计这些更重要的东西。
DolphinScheduler 还有一个很实用的特性:可以拖拽多个任务节点组成 DAG,节点之间能设依赖关系、超时时间和重试次数。这对模拟银行日终批量跑批特别合适。比如先抽客户表,再抽交易流水表,两张表都成功后才触发一个数据校验节点。这种“先 A 后 B,全成功才继续”的编排逻辑,在真实数仓项目里是刚需。所以技术选型不是越重越好,而是看它能不能把你想练的逻辑清晰地暴露出来。
2. 动手前先画清楚:表结构、字段映射与抽取策略
2.1 环境准备清单,别漏掉关键组件
我第一次搭环境时想当然,以为装好 DolphinScheduler 就能直接跑,结果缺了一堆依赖。如果你是从零开始,建议按这份清单准备:
| 组件 | 用途 | 版本建议 |
|---|---|---|
| MySQL 8.x | 源库和目标库 | 8.0+,注意驱动兼容性 |
| DolphinScheduler | 调度和工作流编排 | 3.x,界面和 API 更成熟 |
| JDK | DolphinScheduler 运行依赖 | JDK 8 或 11 |
| ZooKeeper | DolphinScheduler 集群协调 | 3.6+,单机练习也要装 |
| MySQL JDBC 驱动 | 让 DolphinScheduler 能连 MySQL | mysql-connector-java 8.x |
这里特别提醒:DolphinScheduler 3.x 之后,即使单机部署,ZooKeeper 也是启动 Master/Worker 的必要组件。我当时跳过了它,结果服务一直起不来,看日志才发现是注册中心没连上。另外,JDBC 驱动不要只放在服务端 lib 目录,还要确认 Worker 节点执行 SQL 任务时能加载到驱动,否则会出现“数据源连通性测试正常,但任务跑起来报找不到驱动”的诡异问题。
2.2 源表和目标表怎么设计:字段映射与类型转换
小练习的数据源是 MySQL,目标也是 MySQL,看起来类型天然一致,但真正核对字段时还是会发现不少坑。比如源库id_no是加密或者脱敏后的字符串,目标层不希望存原始值;源库 DECIMAL(18,2),目标库如果误建成 FLOAT 会出现精度丢失;源库 DATETIME 带时分秒,目标库如果只想保留到天,就需要在抽取 SQL 里显式转换。
我的做法是先把映射表写清楚,再建目标表。下面是我练习时用的映射关系示例:
| 源表字段 | 源类型 | 目标字段 | 目标类型 | 处理逻辑 |
|---|---|---|---|---|
| cust_id | VARCHAR(32) | cust_id | VARCHAR(32) | 原样复制 |
| cust_name | VARCHAR(128) | cust_name | VARCHAR(128) | 原样复制 |
| id_no | VARCHAR(64) | id_no_masked | VARCHAR(64) | 保留后四位,其余打码 |
| mobile | VARCHAR(20) | mobile | VARCHAR(20) | 原样复制 |
| create_time | DATETIME | etl_time | DATETIME | 改成当前抽取时间 |
有人会觉得这么简单没必要写文档,但相信我,当你要同时维护多张表、多个调度任务时,字段映射表就是你的救命稻草。特别是“哪些字段要做清洗”这栏,直接决定后续数据质量。
目标表建表时,我额外加了两个字段,一个是data_date(业务日期),一个是etl_time(抽取时间)。data_date很重要,因为增量抽取只会拉“某一天”的数据,如果目标表没有这个字段,重跑时你根本不知道这行数据对应哪天的业务,回溯和核对都无从谈起。
2.3 全量抽取还是增量抽取:用哪种,看什么
这是每个做数据抽取的人都会面临的选择。练习里正好有两张差异很大的表,非常适合对比学习。
客户表是典型的缓慢变化维,数据量不大,每天被改动的行数有限,但我练习时故意用了“全量抽取+目标表先清后插”的方式。原因很简单:客户表如果要做增量,必须依赖update_time字段准确且可靠,但这个字段经常会被业务系统漏更新。数据量不大时,全量抽取是成本最低、正确性最高的方案。
交易流水表则必须用增量抽取。因为每天可能新增几十万条数据,全量抽取不仅慢,还会对源库产生很大压力。增量条件我用了txn_time而不是create_time,为什么?txn_time是真实的交易发生时间,create_time是记录插入时间。如果上游系统补录了昨天的交易,用create_time做增量会漏掉补录数据,用txn_time才能保证业务意图。
全量和增量的对比我整理了一个简表:
| 维度 | 全量抽取 | 增量抽取 |
|---|---|---|
| 适用数据量 | 小表、维表 | 大表、流水表 |
| 实现复杂度 | 低,直接 TRUNCATE+INSERT | 高,需维护水位线或时间条件 |
| 对源库压力 | 大,全表扫描 | 小,利用索引范围扫描 |
| 重跑幂等性 | 好,先清后插 | 需要额外处理,防止重复插入 |
| 依赖字段 | 无 | 业务字段或自增 ID |
练习的时候不要只做一个方案。我建议同一张表先用全量跑通,再改成增量,对比两次运行的时间、资源占用和返回值,这样你对“为什么生产环境要费劲做增量”会有直观感受。
3. 在 DolphinScheduler 上搭建抽取工作流
3.1 数据源配置:最常见的问题都出在这一步
DolphinScheduler 里的数据源配置是整个项目的入口,也是我第一次踩坑最多的地方。进入“数据源中心”,新建 MySQL 数据源,需要填数据库地址、端口、数据库名、用户名和密码。看起来很简单,但有几个细节:
第一,jdbc:mysql://地址后面的参数要加useUnicode=true&characterEncoding=utf8&useSSL=false&allowPublicKeyRetrieval=true。否则中文数据抽到目标库容易变问号,MySQL 8 还可能出现连不上报Public Key Retrieval is not allowed的问题。
第二,数据源要区分“源库”和“目标库”,但 DolphinScheduler 里允许同一个数据源被不同的 SQL 节点使用。我建议还是分开建两个数据源,命名带_src和_ods后缀,避免后续看工作流时分不清楚在抽哪个库。这个习惯在生产协作里非常重要。
第三,数据源配置保存前一定要点“测试连接”。如果你在数据库客户端里能连上,但这里测试失败,十有八九是驱动版本不对。DolphinScheduler 3.x 默认自带 MySQL 驱动,但如果版本不匹配,需要在每个 Worker 节点的 lib 目录下替换驱动,然后重启 Worker 进程。我当时遇到一个奇怪现象:Master 节点测试连接成功,但 Worker 跑任务时失败,最后发现就是因为驱动只替换了 Master 的目录。
3.2 创建任务节点:SQL 查询、目标写入和参数
数据源配好后,创建一个工作流,在画布上拖一个 SQL 节点。这里要理解 DolphinScheduler 的 SQL 节点执行逻辑:它会在你选定的数据源上执行一段 SQL,然后把结果集写入到另一个数据源指定的表中。
所以一个典型的抽取任务可以拆成三步:
第一步,在“数据源”下拉框里选择“源库数据源”,输入查询 SQL。比如抽客户表:
SELECT cust_id, cust_name, CONCAT('****', RIGHT(id_no, 4)) AS id_no_masked, mobile, NOW() AS etl_time, '${business_date}' AS data_date FROM customer WHERE update_time >= '${business_date} 00:00:00' AND update_time < DATE_ADD('${business_date}', INTERVAL 1 DAY);第二步,在 SQL 节点的“目标数据源”里选择“ODS 库数据源”,并指定目标表ods_customer。DolphinScheduler 会自动把查询结果插入目标表。这里要注意:如果目标表不存在,SQL 节点不会自动建表,你需要先在目标库里手动建好表结构。
第三步,处理重跑时的重复数据。SQL 节点默认是执行 INSERT,如果目标表已经有同一天的数据,重跑就会造成重复。我的做法是在查询之前先去目标库执行一个 DELETE 语句,把data_date等于当天业务日期的老数据清掉。具体实现是在工作流里再加一个 SQL 节点,数据源选 ODS 库,SQL 写:
DELETE FROM ods_customer WHERE data_date = '${business_date}';这里${business_date}是 DolphinScheduler 的全局参数。我配置了一个日期参数,默认值为昨天,格式yyyy-MM-dd。节点之间用依赖关系串起来,先执行清理节点,再执行抽取节点,最后用字段参数和全局参数一起控制增量范围。
3.3 调度周期与任务依赖:模拟银行日终批量跑批
工作流已经能手动跑通之后,就要加上调度。DolphinScheduler 的调度配置在“工作流定义”页面,点击“定时”按钮,设置 crontab。银行日终跑批一般是凌晨,我设置的是每天 1 点 5 分:
0 5 1 * * ?这个表达式的意思是每天 01:05:00 触发。为什么选 1 点 5 分?因为源库白天还在频繁写入,凌晨数据基本稳定,且要给上游系统留出完成日切的时间。练习阶段你可以改成白天的时间方便观察,但逻辑上要理解“调度时间必须晚于数据就绪时间”。
DolphinScheduler 还支持跨节点依赖和补数。如果我需要跑某一天的历史数据,可以在“工作流实例”页面选择日期,手动补跑。这个功能对银行数据回刷特别重要,一定要试一次。我当时的验证方式是:先把调度时间去掉,只做手动触发;手动跑通后再开定时,观察第二天的调度实例是否准时生成。确认没问题后,再往工作流里加“依赖节点”,让两个抽取任务并发执行,而不是串行等待。
4. 实测后的排错链路:从失败到稳定运行
4.1 驱动、时区、字符集:第一次运行失败的三座山
第一次真正跑调度的时候,我守着屏幕看任务从 “运行中” 变成 “失败”,一时不知道从哪查起。后来总结出一套排查顺序:先看日志,再看参数,最后看数据。
第一座山是驱动问题。报错信息里出现No suitable driver或者Cannot load driver class,基本都是驱动没放对位置。DolphinScheduler 的 Worker 节点执行 SQL 时,会从自己的 lib 目录加载 JDBC 驱动。如果你只在数据源中心测试连接成功,不代表 Worker 节点成功后就能加载到驱动。解决方法是确认所有 Worker 节点的lib目录下都有对应驱动,并重启 Worker 服务。
第二座山是时区问题。MySQL 连接串如果没有带serverTimezone=Asia/Shanghai,DolphinScheduler 执行NOW()或日期比较时,可能出现和本地时间相差 8 小时的情况。我踩到的现象是:明明凌晨 1 点跑的调度,目标表里的etl_time却写成了前一天下午 5 点。后来把连接串里显式加上serverTimezone=Asia/Shanghai才解决。
第三座山是字符集。源库表有中文,SQL 查询结果写入目标库后变成乱码。问题通常出在数据库连接字符集设置。连接串里必须有characterEncoding=utf8。同时要检查目标表、目标库的 charset 是不是utf8mb4,避免某些生僻字或表情符号无法存储。练习时最好从建表阶段就统一字符集,不然后期改起来很麻烦。
4.2 重复数据与漏数据:增量抽取的边界问题
数据能跑通之后,我开始检查数据的正确性,结果发现两个数据质量问题:重复和漏数。
重复的原因是我第一次没有做幂等设计。调度任务第一次跑成功后,我为了验证又手动重跑了一次,结果目标表里出现了两倍的交易流水。解决方法是刚才提到的:先按data_date清理目标表,再插入当天数据。这样无论任务重跑多少次,只要在同一个业务日期下,结果都是一样的。
漏数的问题更隐蔽。源库的交易流水表txn_time上确实有索引,但源库的写入程序有一个批量提交机制,部分数据是在凌晨 0 点 2 分才插入的,而我的调度是 0 点 5 分跑的。我当时用txn_time >= '业务日期 00:00:00'作为增量条件,结果这些批量延迟写入的数据被漏掉了。也就是说,只按业务时间范围抽一次,并不能保证数据全部到位。
一个常见解法是“数据日期偏移”。把增量查询的时间范围往前多扩一点,比如查询头一天晚上的数据:
WHERE txn_time >= '${business_date} 00:00:00' AND txn_time < DATE_ADD('${business_date}', INTERVAL 1 DAY)虽然这个条件没有直接解决延迟写入,但我在生产项目里通常会再配合“定时补偿调度”:白天每隔几小时再抽一次前一天的增量,或者通过核对源表最大txn_time来判断当天数据是否完整。小练习里,更简单的做法是:把调度时间调到凌晨 2 点以后,同时重跑当天数据作为兜底。我最后采用的是“1 点 5 分主调度 + 2 点 30 分补数调度”,两个工作流都执行同样的抽取逻辑,由于目标表有先清后插的幂等设计,重复执行不会产生脏数据。
4.3 调度延迟和失败重试:让数据准点进入数仓
数据质量和调度稳定性是分不开的。DolphinScheduler 中的每个任务都可以单独设置失败重试次数和间隔。我当时给抽取节点配置了“失败重试 3 次,每次间隔 1 分钟”,避免因为源库连接闪断导致整个工作流失败。这里有个容易被忽略的点:重试可能会重复执行同一个任务,所以“先清后插”的幂等逻辑必须写清楚,否则重试一次就会多一份数据。
调度延迟我也遇到过。源库在凌晨有大事务跑批,导致我在凌晨 1 点发起的查询长时间拿不到锁,任务等待了几分钟才执行。这个问题在真实银行环境更明显。练习阶段我没有引入复杂的高可用方案,而是给 SQL 节点设置了查询超时时间。DolphinScheduler 的 SQL 节点有“超时告警”配置,一旦超过设定时间就标记失败并重试。同时我改成了“凌晨 2 点 30 分补跑”的策略,给上游跑批留出足够的窗口。
一个很实用的诊断技巧是:在工作流里加一个“测试节点”,只执行一句SELECT 1,用来验证整个调度链路是否健康。如果连SELECT 1都失败,说明是环境或连接问题,而不是抽取逻辑问题。这个节点看起来不够“高级”,但排错时能省很多时间。
5. 练习之后的进阶方向与个人习惯
5.1 从单表同步到多表依赖:把 DAG 用起来
4 张表都调通后,我开始尝试把它们串成一个有依赖的 DAG。DolphinScheduler 的 DAG 不仅能表达“先抽客户表,再抽流水表”,还能做更复杂的编排。我设计的结构是:
- 任务 A:清理 ODS 层客户表当天数据。
- 任务 B:清理 ODS 层流水表当天数据。
- 任务 C:抽取客户表数据。
- 任务 D:抽取交易流水表数据。
- 任务 E:数据校验与通知。
A 和 B 可以并行,C 依赖 A,D 依赖 B,E 依赖 C 和 D 同时成功。这样如果某一张源表挂了,不会影响另一张表的抽取;只有两张表都成功,才会进入校验节点。校验节点可以做简单的行数比对:查询源表和目标表的记录数,如果不一致就把工作流置为失败。这就是 DAG 的魅力,它把“并行”“依赖”“全局决策”都可视化出来。
5.2 加上数据质量校验:行数比对、空值检查和重跑机制
数据抽取不是“抽完就完事”,还得确认数据没丢、没多、没脏。我在练习时加了三个校验点:第一个是源表和目标表的总行数比对,适合全量抽取的客户表;第二个是交易流水表的增量条数校验,我会用 SQL 统计源表当天符合条件的记录数,再和目标表当天记录数做对比;第三个是空值检查,重点看流水金额是否为空、客户证件号打码后是否变成纯空字符串。
行数比对的 SQL 示例:
-- 源表当天流水数 SELECT COUNT(*) AS src_cnt FROM txn_log WHERE txn_time >= '${business_date} 00:00:00' AND txn_time < DATE_ADD('${business_date}', INTERVAL 1 DAY);把上面这个查询的结果和目标表ods_txn_log中data_date等于当天日期的条数做对比。如果两边不一致,就发告警。我个人的习惯是,校验节点不要用 DolphinScheduler 的默认告警,而是在校验失败时抛出一个异常,让整个工作流标记为失败,并在告警信息里附上源表条数和目标表条数的差值。这样第二天查看任务历史时,一眼就能定位问题。
5.3 我给初学者的一些实操习惯
最后分享几个我这次练习里形成的习惯,不保证最优,但确实能减少很多折腾。第一个习惯是“先手动,后定时”。任何新加的表或修改过的抽取逻辑,都先手动触发一次,确认数据和日志都没问题后,再挂到调度上。不要直接改定时配置,否则出了问题都不知道是调度问题还是逻辑问题。
第二个习惯是“把业务日期作为参数引出来”。不要在每个节点里硬编码日期,而是定义一个全局参数business_date,默认用 DolphinScheduler 的$[yyyy-MM-dd]自动取前一天。这样补数时只需要修改参数值,不用改每个 SQL。
第三个习惯是“注意目标表的主键和唯一键”。我练习时因为目标表没建唯一索引,导致重复数据在 SQL 层无法防住。最后我在ods_txn_log表上加了联合唯一键(txn_id, data_date),这样即使 INSERT 语句因为重试重复执行,数据库层也会直接拒绝重复数据,等于多了一道保险。对于源表没有明确主键的字段,也要先通过分组检查确认哪些字段组合能唯一标识一行,再决定目标表的约束。
这个小练习里我最大的体会是:数据抽取看着是“复制粘贴”,其实每一步都在回答“数据怎么来、怎么存、怎么算”。把客户表和交易流水表完整地跑通一个闭环之后,你再去看 DataX、SeaTunnel 这些工具,思路会清楚很多。下次我打算在这个基础上把 ODS 到 DWD 的清洗转换也加进去,让整个数仓链路更完整。