ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

基于DolphinScheduler的银行数据抽取全量与增量实践

2026/9/8 2:53:48 拓冰建站 浏览量
基于DolphinScheduler的银行数据抽取全量与增量实践 前段时间我拿一个银行场景练了练数据抽取任务本身不复杂从几个业务库里把客户表和交易流水表抽到数仓的 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 更成熟JDKDolphinScheduler 运行依赖JDK 8 或 11ZooKeeperDolphinScheduler 集群协调3.6单机练习也要装MySQL JDBC 驱动让 DolphinScheduler 能连 MySQLmysql-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_idVARCHAR(32)cust_idVARCHAR(32)原样复制cust_nameVARCHAR(128)cust_nameVARCHAR(128)原样复制id_noVARCHAR(64)id_no_maskedVARCHAR(64)保留后四位其余打码mobileVARCHAR(20)mobileVARCHAR(20)原样复制create_timeDATETIMEetl_timeDATETIME改成当前抽取时间有人会觉得这么简单没必要写文档但相信我当你要同时维护多张表、多个调度任务时字段映射表就是你的救命稻草。特别是“哪些字段要做清洗”这栏直接决定后续数据质量。目标表建表时我额外加了两个字段一个是data_date业务日期一个是etl_time抽取时间。data_date很重要因为增量抽取只会拉“某一天”的数据如果目标表没有这个字段重跑时你根本不知道这行数据对应哪天的业务回溯和核对都无从谈起。2.3 全量抽取还是增量抽取用哪种看什么这是每个做数据抽取的人都会面临的选择。练习里正好有两张差异很大的表非常适合对比学习。客户表是典型的缓慢变化维数据量不大每天被改动的行数有限但我练习时故意用了“全量抽取目标表先清后插”的方式。原因很简单客户表如果要做增量必须依赖update_time字段准确且可靠但这个字段经常会被业务系统漏更新。数据量不大时全量抽取是成本最低、正确性最高的方案。交易流水表则必须用增量抽取。因为每天可能新增几十万条数据全量抽取不仅慢还会对源库产生很大压力。增量条件我用了txn_time而不是create_time为什么txn_time是真实的交易发生时间create_time是记录插入时间。如果上游系统补录了昨天的交易用create_time做增量会漏掉补录数据用txn_time才能保证业务意图。全量和增量的对比我整理了一个简表维度全量抽取增量抽取适用数据量小表、维表大表、流水表实现复杂度低直接 TRUNCATEINSERT高需维护水位线或时间条件对源库压力大全表扫描小利用索引范围扫描重跑幂等性好先清后插需要额外处理防止重复插入依赖字段无业务字段或自增 ID练习的时候不要只做一个方案。我建议同一张表先用全量跑通再改成增量对比两次运行的时间、资源占用和返回值这样你对“为什么生产环境要费劲做增量”会有直观感受。3. 在 DolphinScheduler 上搭建抽取工作流3.1 数据源配置最常见的问题都出在这一步DolphinScheduler 里的数据源配置是整个项目的入口也是我第一次踩坑最多的地方。进入“数据源中心”新建 MySQL 数据源需要填数据库地址、端口、数据库名、用户名和密码。看起来很简单但有几个细节第一jdbc:mysql://地址后面的参数要加useUnicodetruecharacterEncodingutf8useSSLfalseallowPublicKeyRetrievaltrue。否则中文数据抽到目标库容易变问号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 连接串如果没有带serverTimezoneAsia/ShanghaiDolphinScheduler 执行NOW()或日期比较时可能出现和本地时间相差 8 小时的情况。我踩到的现象是明明凌晨 1 点跑的调度目标表里的etl_time却写成了前一天下午 5 点。后来把连接串里显式加上serverTimezoneAsia/Shanghai才解决。第三座山是字符集。源库表有中文SQL 查询结果写入目标库后变成乱码。问题通常出在数据库连接字符集设置。连接串里必须有characterEncodingutf8。同时要检查目标表、目标库的 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 依赖 AD 依赖 BE 依赖 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 的清洗转换也加进去让整个数仓链路更完整。