背景
很多公司的数据集成,都是从一堆脚本开始的。
最早是几个 crontab 定时任务,用 mysqldump 把业务库导到数仓;后来数据源多了,又加了一堆 Python 脚本、Shell 脚本、Sqoop 命令;再后来,脚本散落在各个服务器上,没人知道哪个脚本在跑、跑没跑成功、数据对不对。
等到你意识到"这样不行"的时候,数据集成已经变成了一团乱麻——脚本上百个、服务器十几台、依赖关系没人理得清、出了问题全靠人肉排查。
这篇文章记录数据集成平台从"脚本"到"平台化"的完整演进过程,以及每个阶段要解决的核心问题。
一、脚本时代:能跑就行,但迟早会崩
脚本时代的特点,一句话概括:能用,但不可维护。
典型场景是这样的:
bash
# crontab 里的一个数据同步任务0 2 * * * /home/data/scripts/sync_order.sh >> /home/data/logs/sync_order.log 2>&1
sync_order.sh 里是一堆 sqoop import 命令:
bash
#!/bin/bashsqoop import \ --connect jdbc:mysql://db1:3306/order_db \ --username data_etl \ --password xxx \ --table orders \ --target-dir /warehouse/ods/orders \ --incremental append \ --check-column order_id \ --last-value $(cat /home/data/last_value/orders)
脚本时代的痛点,我列几个最典型的:
- 脚本散落,无人知晓:脚本在 A 服务器、B 服务器、C 服务器上各有一堆,没人知道总共多少个、都在干什么。
- 依赖关系混乱:任务 A 依赖任务 B 的结果,但依赖关系写在脚本的注释里,甚至写在开发者的脑子里。
- 失败无感知:任务跑失败了,除非有人去看日志,否则没人知道。数据断了,下游报表错了,才发现。
- 密码明文:数据库密码直接写在脚本里,安全风险极高。
- 无法复用:每个脚本都是孤立的,同样的同步逻辑复制粘贴了十遍。
脚本时代的核心矛盾:数据集成在"野蛮生长",但缺乏统一的管控和治理。
二、平台化要解决的核心问题
从脚本到平台,不是"把脚本搬到一个系统里",而是要解决几个根本性的问题。
2.1 任务统一管理
所有数据集成任务,统一在一个平台上创建、配置、调度、监控。不再有散落的脚本,不再有"不知道在哪跑"的任务。
2.2 依赖关系可视化
任务之间的依赖关系(A 完成后才能跑 B),用 DAG(有向无环图)表达,可视化展示,而不是写在注释里。
2.3 失败告警与重试
任务失败自动告警(钉钉/邮件),支持自动重试,支持断点续传。
2.4 数据源统一接入
数据库、文件、API、消息队列等各种数据源,统一接入管理,连接信息集中配置,不再散落在脚本里。
2.5 任务可复用
把"从 MySQL 同步到 Hive"这种通用能力抽象成"同步任务模板",配置一次,复用多次。
三、平台架构设计
数据集成平台的典型架构,分四层:
┌─────────────────────────────────────────────────────┐│ 接入层 │ │ Web 管理界面 │ REST API │ 开放平台(任务提交) │ ├─────────────────────────────────────────────────────┤ │ 服务层 │ │ 任务管理 │ 调度引擎 │ 监控告警 │ 数据源管理 │ ├─────────────────────────────────────────────────────┤ │ 执行层 │ │ 同步引擎(Reader/Writer) │ 执行器集群 │ CDC 引擎 │ ├─────────────────────────────────────────────────────┤ │ 存储层 │ │ 任务元数据(MySQL) │ 运行日志(ES) │ 监控指标 │ └─────────────────────────────────────────────────────┘
3.1 核心概念:Reader / Writer 插件化
数据集成平台最核心的设计,是 Reader / Writer 插件化。这是 DataX 的经典设计思想,也是几乎所有数据集成平台的通用范式。
- Reader:从数据源读数据的插件(MySQL Reader、Oracle Reader、HDFS Reader、Kafka Reader...)
- Writer:往目标写数据的插件(Hive Writer、ClickHouse Writer、ES Writer...)
一个同步任务 = 一个 Reader + 一个 Writer。Reader 和 Writer 可以任意组合,实现"任意源到任意目标"的同步。
MySQL Reader ──┐Oracle Reader ──┤ HDFS Reader ───┼──→ 同步引擎 ──→ Hive Writer Kafka Reader ──┤ ├─→ ClickHouse Writer API Reader ──┘ └─→ ES Writer
插件化的好处:
- 扩展性强:新增一个数据源,只需开发一个 Reader/Writer 插件,不用改核心引擎。
- 组合灵活:M×N 种同步路径,只需要 M+N 个插件,而不是 M×N 套代码。
- 职责单一:每个插件只负责一种数据源的读写,代码清晰可维护。
3.2 调度引擎设计
调度引擎是平台的"大脑",负责决定"什么任务、什么时候跑、按什么顺序跑"。
核心能力:
- 定时调度:Cron 表达式,支持分钟/小时/天/周/月各种粒度
- 依赖调度:DAG 依赖,上游任务完成后触发下游
- 并发控制:限制同时运行的任务数,防止资源打满
- 优先级:重要任务优先调度
调度引擎的实现,自研和开源两条路:
| 方案 | 优点 | 缺点 |
| 自研调度 | 完全可控,可深度定制 | 开发成本高,边界情况多 |
| DolphinScheduler | 成熟稳定,DAG 可视化,社区活跃 | 定制化受限 |
| Airflow | 生态丰富,Python 友好 | 重,学习成本高,调度延迟 |
| Azkaban | 简单易用 | 功能相对少,社区不活跃 |
我们最终选了 DolphinScheduler 作为调度底座,在其上做二次开发。原因:DAG 可视化做得好、支持工作流编排、社区活跃、国产开源(文档和社区都是中文友好)。
3.3 执行器集群设计
同步任务的执行,需要分布式执行器集群,支持水平扩展。
调度引擎 → 下发任务到执行器 → 执行器(Worker 节点)执行同步 → 上报执行状态和日志
执行器集群的关键设计:
- 水平扩展:任务量大时加 Worker 节点即可
- 资源隔离:不同优先级的任务分配到不同的 Worker 队列
- 故障转移:Worker 挂了,任务自动转移到其他 Worker
3.4 监控告警设计
监控告警是平台"能不能放心用"的关键。
监控三个维度:
- 任务维度:任务成功率、失败率、平均耗时、延迟情况
- 数据维度:同步的数据量、数据条数、数据质量(空值率、重复率)
- 资源维度:Worker 的 CPU、内存、网络、磁盘
告警策略:
- 任务失败:立即告警(钉钉 + 邮件)
- 任务延迟:超过阈值告警
- 数据量异常:同步的数据量骤降/骤增告警(可能是源端数据异常)
- 数据质量异常:空值率、重复率超过阈值告警
四、CDC:从批量同步到实时同步
传统的数据集成是"批量同步"——T+1 定时全量或增量同步。但业务对实时性的要求越来越高,CDC(Change Data Capture,变更数据捕获)成了数据集成平台的标配能力。
4.1 CDC 的原理
CDC 的核心思想:监听数据库的变更日志(Binlog / Redo Log),实时捕获数据的增删改,同步到下游。
MySQL Binlog → Canal / Flink CDC → Kafka → 下游(数仓/缓存/索引)
4.2 CDC 方案选型
| 方案 | 原理 | 优点 | 缺点 |
| Canal | 伪装成 MySQL Slave,解析 Binlog | 成熟稳定,阿里开源 | 只支持 MySQL,需部署 Canal Server |
| Flink CDC | Flink 内置 CDC 连接器 | 流批一体,无中间组件 | 依赖 Flink,版本兼容性需注意 |
| Debezium | Kafka Connect 插件 | 支持多种数据库,生态好 | 依赖 Kafka Connect,部署复杂 |
我们选了 Flink CDC,理由是:团队已经有 Flink 集群,直接复用,不用额外部署 Canal Server;而且 Flink CDC 支持流批一体,批量同步和实时同步可以统一在 Flink 上做。
4.3 CDC 的坑
CDC 不是银弹,有几个坑要提前知道:
- 大事务导致延迟:源端一个超大事务(比如一次更新 1000 万行),Binlog 解析和下游写入会卡很久。
- DDL 变更处理:源端表结构变更(加字段、改类型),CDC 链路需要正确处理,否则数据同步会断。
- 全量 + 增量衔接:第一次同步需要全量 + 增量的无缝衔接,中间不能丢数据。
五、演进路径总结
数据集成平台的演进,我总结成四个阶段:
| 阶段 | 形态 | 核心问题 | 关键动作 |
| 脚本时代 | crontab + 脚本 | 不可维护、不可监控 | 无 |
| 任务调度化 | 调度系统 + 脚本 | 任务可管理,但同步逻辑还是脚本 | 引入调度系统 |
| 平台化 | 集成平台 + 插件化 | 同步能力可复用、可扩展 | Reader/Writer 插件化 |
| 实时化 | 平台 + CDC | 批量同步无法满足实时需求 | 引入 CDC |
每个阶段的跃迁,都是由"业务需求倒逼"的。 数据量大了 → 脚本跑不动了 → 上调度系统;数据源多了 → 脚本写不过来了 → 上插件化平台;实时需求来了 → 批量同步不够了 → 上 CDC。
六、最后说一句
数据集成平台,看起来是个"搬数据"的工具,技术含量不高。但真正做起来,你会发现它涉及调度、分布式、插件化、CDC、监控告警一大堆工程问题。
从脚本到平台,本质是从"人肉运维"到"工程化"的跃迁。 脚本时代靠的是人的经验和责任心,平台时代靠的是系统的能力和机制。
如果你还在用脚本做数据集成,别觉得"能用就行"。等到数据源多了、任务多了、出问题排查不动了,再想平台化就晚了。
平台化这件事,越早做,成本越低。
本文基于作者团队数据集成平台的建设实践。技术栈参考 DataX + DolphinScheduler + Flink CDC,其他技术栈可类比。