1. 增量数据层(Delta Layer)到底解决什么问题
1.1 从全量同步到增量同步的演进
先说个我经常在团队里听到的问题:为什么已经有了数据仓库,还要单独搞一个 Delta Layer?这玩意儿在传统的数仓分层里(ODS、DWD、DWS、ADS)好像没有直接对应的位置。我的理解是,Delta Layer 本质上不是一个独立的分层,而是一种数据流转策略——它描述的是从业务系统到大数据平台之间,那一层专门承载"增量变化"的数据管道。
在早期做数据同步的时候,大家最爱干的事就是每天凌晨跑一个全量抽取,把业务库的表整个拉一遍到 Hive 里。这种方式在数据量小、业务逻辑简单的年代确实够用。但后来就有两个问题开始冒头:第一是数据量越来越大,全量抽取的时间窗口越来越紧张,凌晨两三个小时的调度窗口根本不够用;第二是业务对数据时效性的要求从 T+1 一路提到了分钟级甚至秒级,全量同步这种"拿昨天的快照做今天的分析"的模式就显得非常尴尬。
所以增量同步就变成了刚需,而 Delta Layer 正是为了承接这种增量数据而设计的。你可以把它理解为数据进入数仓前的一个"缓冲带"——增量数据从业务系统产生后,先落到这一层,经过清洗、格式化、去重,再把干净的结果推给下游。这样一来,下游的计算引擎永远不需要去直接面对原始业务库,也不用反复全量拉取数据,整体链路的压力和成本都会小很多。
从我个人经验看,如果你还停留在全量同步的阶段,不用急着否定这种方案,因为数据量小的时候全量同步反而是最简单、最少出错的。但只要你开始接触 CDC、实时数仓、数据湖这类概念,Delta Layer 就是一个绕不开的桥头堡。
1.2 Delta Layer 在三层数仓架构里的真实位置
传统的数仓分层是 ODS(操作数据存储)、DWD(明细数据层)、DWS(汇总数据层)、ADS(应用数据层)。很多人会把 Delta Layer 直接等同于 ODS,但我个人认为这种理解太粗暴了。ODS 承载的是"原样落地"的原始数据,全量、增量都往里塞,而 Delta Layer 更侧重于"增量变化"的捕获和流转,两者有交集,但设计逻辑并不相同。
我更愿意把 Delta Layer 理解成一个独立于传统分层的技术中间件。在你的业务数据库和数据平台之间,它负责做三件事:一是捕获源端的变化事件(增删改),二是把变化事件转换为统一格式的数据记录(比如 Append 到 Kafka Topic 或写入 Delta Table),三是保证这一层的数据可以被下游安全、高效地消费。
举个实际例子。在某个电商场景里,订单表每分钟有上千条新纪录产生,还有数百条状态更新(比如待付款变成已付款)。如果你用全量同步,每分钟拉一次全表,数据库压力大,下游处理慢。而用 Delta Layer 的思路,每次只把变化的那几百条数据同步过来,效率可以说是数量级的提升。更重要的是,配合 upsert 语义,你可以让下游表的状态始终保持和源库一致,这才是 Delta Layer 真正的核心竞争力。
一句话总结:Delta Layer 不是某个具体的软件,而是一种数据架构策略,它的核心目标就是用最小的 IO 开销,把源端的变化精准地传达到数据平台的每一层。
2. 增量数据的捕获机制
2.1 基于时间戳和版本号的增量抽取
聊完定位,接下来就要落地了。增量捕获的技术手段有很多,我从最轻量级的开始讲。
第一种是时间戳增量抽取。操作非常简单:在源表里加一个update_time字段,每次同步的时候,把where update_time > 上次同步时间点的数据拉出来。这个方案的优点是侵入性极小,不会给数据库带来额外的性能负担,开发起来也很快。缺点是它有个天然的盲区——如果某条记录在上次同步之后被删除了,你是永远发现不了的,因为被删的行不存在了,时间戳查询根本看不出来。
第二种是基于版本号的增量抽取。在一些业务表里,会有一个自增的版本字段或者主键 ID,同步时记录下当前的最大版本号,下次同步拉取比这个版本号大的数据。这种方案相比时间戳更稳定,因为版本号是单调递增的,不会因为跨天、时区的问题导致漏数据。但它同样解决不了数据删除的问题,而且如果业务代码不规范,版本号出现回退,就会导致数据错乱。
我自己的经验是,时间戳和版本号方案适合作为兜底手段,用于那些数据量不大、删除敏感度不高的业务表。比如配置表、品类表、地区表这类低频变化的维度数据,用时间戳增量就完全够用。但一旦涉及订单、库存、余额这种强一致性的核心数据,就必须引入更可靠的机制。
2.2 解析数据库日志的 CDC 方案
真正可靠的增量捕获方案,是解析数据库的 binlog 或者 WAL 日志,这个技术方向叫 CDC(Change Data Capture,变更数据捕获)。CDC 的核心理念是:不主动去查询源表,而是监听数据库的日志,把每一次 insert、update、delete 操作都解析成一条流式事件,再发送给下游。
以最常见的 MySQL 为例,binlog 里记录了每一行数据的变更。Canal、Debezium、Flink CDC 这些工具,本质都是 binlog 的消费者。它们伪装成 MySQL 的从节点,向主库请求 binlog,然后解析出结构化的事件。
这里有一个很大的优势:所有变更事件不会丢失,包括删除操作。你可以完整拿到业务系统的每一步数据变化,真正做到"源库发生了什么,目标端就重放什么"。而且 CDC 对源库的压力非常小,因为它走的是日志复制通道,不会反复查询业务表。
CDC 方案的代价是什么呢?第一是运维复杂度上来了,你要额外部署一个 CDC 组件,还要维护和数据库之间的连接,处理 binlog 过期、连接中断、位点持久化这些技术细节;第二是数据格式的解析有门槛,binlog 是二进制格式,不同版本的 MySQL 解析规则也有差异,排起错来不是那么轻松。但从生产环境的可靠性来看,这套投入是值得的。
2.3 各类增量方案怎么选
我结合自己的使用经验,把常见的增量捕获方案放在一起对比了一下,方便你按需选择:
| 方案 | 侵入性 | 实时性 | 删除感知 | 源库压力 | 运维成本 | 适用场景 |
|---|---|---|---|---|---|---|
| 时间戳增量 | 低 | 分钟级到小时级 | 无法感知 | 低 | 极低 | 低频维度表、配置表 |
| 版本号增量 | 低 | 分钟级到小时级 | 无法感知 | 低 | 极低 | 只追加不修改的日志型表 |
| CDC(binlog/WAL) | 低 | 秒级 | 完整感知 | 极低 | 中高 | 订单、库存、余额等核心业务表 |
| 消息队列埋点 | 高 | 秒级 | 完整感知 | 中 | 中 | 日志系统、事件驱动架构 |
这里说个我在项目里踩过的坑。之前有个系统想保证核心业务数据的实时性,一开始图省事用了时间戳增量,结果运营反馈"退款订单数对不上"。排查下来发现,很多退款订单是先创建、后修改状态,修改的时候时间戳确实变了,但如果某个订单在同步窗口内既被更新又被删除,时间戳方案就完全没法处理。最后换成了 Debezium 做 CDC,才算把这个问题根治掉。
所以我的建议是,核心交易类数据直接上 CDC,省得后面返工;非核心的维度数据,用时间戳或版本号去应付就行,别把所有表都搬到 CDC 上面,不然维护的链路太多,出事的时候排查起来很痛苦。
3. 实操:搭一个生产级的 Delta Layer
3.1 技术选型和完整链路
理论讲再多,不如把链路搭起来跑一遍。我给一套可以直接参考的落地方案:使用 MySQL 作为源业务库,Flink CDC 负责捕获和传输增量,Kafka 作为消息缓冲层,最终落地到 Delta Table(以 Delta Lake 格式存储,底层是 Parquet 文件,挂在 S3 或 HDFS 上)。
链路大致是这样的:
MySQL binlog -> Debezium / Flink CDC -> Kafka Topic -> Flink 流式计算 -> Delta Table -> 下游 Doris / StarRocks / ClickHouse / Hive为什么选 Flink CDC 而不是单独部署 Debezium?核心原因是 Flink CDC 直接和 Flink 的计算引擎打通了,不需要把 Debezium 的变更事件先写进 Kafka 再通过 Kafka Connector 接出来,链路更短,端到端的延迟能做到更低。当然,如果你的技术栈里没有 Flink,用 Debezium + Kafka Connect 也非常成熟,只是多一跳。
关于 Delta Table 的选择,这里要区分一个概念。"Delta Layer" 里的 Delta 是"增量"的意思,而 Delta Lake 是 Databricks 开源的一种数据湖存储格式。两者之间是有联系的——Delta Lake 提供了 ACID 事务、时间旅行、upsert 能力,非常适合作为增量数据落地层的存储引擎。我们把增量数据写进 Delta Table,下游可以直接用 SparkSQL 读取,也可以借助 CDC 的 upsert 语义把数据变更合并到主表。
3.2 环境准备
动手之前先列一下要用的组件版本,参考组合如下:
| 组件 | 版本 | 说明 |
|---|---|---|
| MySQL | 5.7+ | 需要开启 binlog,binlog_format=ROW,binlog_row_image=FULL |
| Flink | 1.17+ | 推荐用 Flink 1.17 以上,flink-cdc 兼容性更好 |
| flink-cdc-connector | 2.4+ | 提供 MySQL CDC 连接器 |
| Delta Lake | 2.4+ | 支持 Flink 写入需要对应版本适配 |
| Kafka | 2.8+ | 作为缓冲层,按需配置 |
关于 MySQL 的 binlog 配置,有几个非常关键的参数,直接决定了 CDC 能不能正常工作:
[mysqld] server-id = 1 log-bin = mysql-bin binlog_format = ROW binlog_row_image = FULL expire_logs_days = 7这里逐条解释一下为什么这么配。binlog_format=ROW是必须的,因为只有行级日志才能记录每一行变更的完整前后镜像;如果是 STATEMENT 格式,记录的是 SQL 语句,解析出来的就不是精确的行变更了。binlog_row_image=FULL也很重要,它的意思是 binlog 里要包含变更行的所有列,而不是只记录修改的列,这样才能保证下游拿到的数据是完整的。expire_logs_days=7是给自己留 7 天的回溯空间,如果 Flink 任务挂了超过 7 天没恢复,binlog 被清理了就只能做全量初始化了。
这里有个基础但容易忘的点:MySQL 的 binlog 默认可能是关闭的,如果你在配置里忘了写log-bin,后面 Flink CDC 任务就直接报错。我见过不止一次有人任务跑不起来,排查了半天才发现 binlog 压根没开。
3.3 Flink CDC 同步任务的配置细节
环境准备好之后,我们来写一个标准的数据同步任务。假设有一张订单表t_order,我们要把它的增量数据同步到 Delta Table,并且以order_id作为主键做 upsert。
用 Flink SQL 的方式相对简单,先建一个 CDC Source 表:
CREATE TABLE source_order ( order_id STRING, user_id BIGINT, order_amount DECIMAL(10, 2), order_status INT, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '192.168.1.100', 'port' = '3306', 'username' = 'cdc_user', 'password' = 'your_password', 'database-name' = 'business_db', 'table-name' = 't_order', 'scan.startup.mode' = 'initial', 'server-id' = '5400-5410' );这段 SQL 里有三个细节值得展开讲。
第一个是scan.startup.mode。我选的是initial,意思是任务启动时会先做一次全量快照,然后自动无缝切换到增量读取。这样做的好处是表里的历史数据不会被遗漏。如果你确认不需要历史数据,只要从当前时刻开始的增量,可以改成latest-offset。但在生产环境,我一般建议用initial,因为后续如果要回溯历史或者做数据修复,有这个快照在会好办很多。
第二个是server-id。Flink CDC 在读取 binlog 时会模拟 MySQL 的从节点,所以需要提供唯一的 server-id。如果同一台机器上跑了多个并发,建议给一个区间,比如5400-5410,这样 Flink 会自动分配,避免不同任务用了同一个 server-id 导致在 MySQL 上互相干扰。
第三个是主键声明。CDC Source 表的PRIMARY KEY是逻辑主键,它不参与 DDL 建表,而是用来描述这条流里"以哪个字段为唯一标识"。下游要做 upsert 的时候,这个主键会直接决定合并逻辑。
接下来是目标表,这里用 Delta Lake 的 Flink Connector 来定义。注意 Delta Lake 官方现在主要支持 Spark 和 Flink 的互动,用 Flink 写 Delta Table 需要引入对应的依赖,并指定一系列连接器参数:
CREATE TABLE delta_order ( order_id STRING, user_id BIGINT, order_amount DECIMAL(10, 2), order_status INT, create_time TIMESTAMP(3), update_time TIMESTAMP(3) ) WITH ( 'connector' = 'delta', 'table.path' = 's3://my-bucket/delta-tables/order', 'delta.appendOnly' = 'false', 'delta.columnMapping.mode' = 'name' );建好两张表之后,一个简单的同步任务就是一行 INSERT INTO:
INSERT INTO delta_order SELECT order_id, user_id, order_amount, order_status, create_time, update_time FROM source_order;但这个只解决了"把数据抄过去"的问题,并没有做真正的 upsert。也就是说,如果源库订单状态从"待付款"改成"已付款",delta_order 里会同时存在两条记录。要处理这种变更,有两种方式。
第一种方式,在 Flink SQL 里用聚合加ROW_NUMBER()去保留最新的一条。比如假设update_time自增:
INSERT INTO delta_order SELECT order_id, user_id, order_amount, order_status, create_time, update_time FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY update_time DESC) AS rn FROM source_order ) t WHERE rn = 1;第二种方式,如果下游是 Delta Lake,我们可以在写完增量之后,用 Spark 对 Delta Table 执行MERGE INTO做合并:
MERGE INTO delta_order AS t USING source_order AS s ON t.order_id = s.order_id WHEN MATCHED THEN UPDATE SET t.user_id = s.user_id, t.order_amount = s.order_amount, t.order_status = s.order_status, t.update_time = s.update_time WHEN NOT MATCHED THEN INSERT (order_id, user_id, order_amount, order_status, create_time, update_time) VALUES (s.order_id, s.user_id, s.order_amount, s.order_status, s.create_time, s.update_time);3.4 增量表的设计和分区策略
增量层的表设计和传统数仓的表设计有一个显著区别:传统数仓分区一般按天、按月,而增量层的数据时效性短,很多场景下需要按小时甚至直接不分区分区。
我在做增量表的时候,一般会遵循几个原则。第一,保留源数据的所有业务字段,不要因为在"增量层"就随便裁剪字段,因为你不知道下游什么时候会需要这些字段。第二,必须带上op_ts(操作时间戳)和record_type(操作类型,insert/update/delete)这类元数据字段,方便下游判断这条数据的语义。第三,按天分区,在分区目录里存放当天所有的变更记录,既方便查询,也方便按天清理。
简单给出一个表结构的参考:
CREATE TABLE delta_ods_order ( order_id STRING, user_id BIGINT, order_amount DECIMAL(10, 2), order_status INT, create_time TIMESTAMP(3), update_time TIMESTAMP(3), op_ts TIMESTAMP(3), event_type STRING, partition_date STRING ) USING DELTA PARTITIONED BY (partition_date);关于表结构这里多说一句,event_type用来标记这条记录是 insert、update 还是 delete。虽然在 CDC 语义里,delete 事件也会被传出来,但如果你不加这个字段,下游就不知道这条数据在源端已经没了,最后做数据核对的时候会对不上。加上这个字段之后,下游可做过滤,也可做"软删除"处理。
分区的粒度,我之前踩过一个坑。有一回去做实时订单分析,下游团队直接查 Hive 上同步过来的 Delta 表,因为分区粒度太大了,查询的时候扫描了整个分区的数据,导致响应很慢。后来我把分区粒度从"天"改成"小时",情况改善了很多。当然,分区粒度越细,文件数量也会越多,小文件问题就跟着来了。这个后面在问题排查部分专门展开。
4. 实战中躲不开的五个坑
4.1 小文件太多,性能反而更差
增量数据有一个天然的副作用,就是写入频率高、单次数据量小。如果你用 Flink 流式写入 Delta Lake,默认几分钟一个 checkpoint,每个 checkpoint 产生一批小文件,一天下来就会出现几千个甚至几万个大小只有几 MB 的文件。下游跑 SQL 的时候,每次都要打开大量小文件做数据扫描,性能反而不如全量同步时的那种大文件查询。
解决办法有几种。第一种是开启 Delta Lake 的自动优化,在写入时增加文件合并策略,比如设置delta.autoOptimize.optimizeWrite = true和delta.autoOptimize.autoCompact = true,这样在写入过程中就会自动将小文件合并成更大的文件。第二种是定时做OPTIMIZE,通过任务在低峰期把 Delta Table 的文件整理一遍。第三种,也是最推荐的,就是在设计分区粒度时把写入频率和文件大小结合起来考虑,比如每 10 分钟触发一次文件合并,保证单个文件不低于 128MB。
这个坑的典型表现是:数据同步任务完全正常,但下游报表查询突然变慢,一看监控,扫描的文件数量是之前的几十倍。如果遇到这种情况,先不要调计算引擎的参数,先看看是不是小文件太多。
4.2 Exactly-Once 语义下的重复数据和数据丢失
这里要澄清一个常见的误区:Flink CDC 本身默认是 At-Least-Once,不是 Exactly-Once。如果 Flink 任务重启,binlog 的位点回退了,那源端的一部分变更就会被重复消费。重复消费本身不可怕,可怕的是下游不做去重,就导致数据翻倍。
我处理这种情况的方式是在 Delta Table 的写入逻辑上保证幂等。具体来说,就是按照主键做 upsert,而不是简单 append。你可以在 Flink 端用ROW_NUMBER() OVER (PARTITION BY 主键 ORDER BY 操作时间戳 DESC)去重,也可以在下游靠 Delta Lake 的MERGE INTO来保证同一主键只有一条记录。更严谨一点,建议把任务配置成 Flink 的 Checkpoint 模式,在故障恢复时从最近一次 checkpoint 的位点开始消费,减少重复窗口。
反过来,也可能出现数据丢失。我遇到过一种情况是 Flink 的 checkpoint 长时间没有完成,而外部系统在 checkpoint 之后把状态清掉了,恢复的时候发现丢了一段数据。这种问题排查起来比较费劲,我的建议是按小时做一次数据比对(源库主键集合 vs 目标表主键集合),及时发现差异,再通过全量初始化来修复。
4.3 数据乱序和时间字段的陷阱
CDC 事件来自多个并发事务,到达目标端的时间顺序并不能完全代表业务上真实的发生顺序。比如订单先创建后被取消,如果取消事件先被消费,创建事件后被消费,如果不做处理,目标表最终保留的是"创建"状态,数据就错了。
解决方案是在确定主键和排序字段时,不能用到达时间,而要使用业务时间或者 binlog 里的操作时间戳。在 Flink SQL 里可以给表声明 Watermark 和 Event Time,再配合ROW_NUMBER()按业务时间排序,确保最终保留的是业务上最新的一条。还有一个小细节:MySQL binlog 的DML事件里有时区信息,如果你的任务没有做时区对齐,很可能出现跨时区的时间偏移。统一用 Asia/Shanghai 时区起步,能少很多麻烦。
4.4 Schema 变更怎么同步
源端加了一个字段,这是做增量同步时一定会遇到的事。Flink CDC 可以感知到 DDL 变更,但同步到 Delta Table 的处理策略要想清楚。
我推荐的方式是:在 Delta Table 的写入任务里开启allowColumnDefaults这类列的自动扩展能力(不同引擎叫法不同)。它的逻辑是当源端新增列时,自动在目标表上以空值或默认值来补齐该列。这个方案实现成本低,适合大多数场景。如果你要求完全同步 DDL,那就需要用 Debezium 的 Schema Evolution 或者 Flink CDC 的 Schema 变更路由,把 DDL 事件单独发给一个管理任务,由它去执行目标端的 DDL。这个方案更重,但对严格一致性的场景(比如金融数据)是必要的。
我个人经验是,中小规模的业务表,直接让 Delta 表自动加列就够了,别为了追求完美去做完整的 Schema 同步,维护成本太高,而且容易出问题。
4.5 消费延迟是普遍问题
增量同步链路搭好以后,最常收到的一个报警就是"数据延迟"。很多人第一时间去查 Flink 任务的吞吐量,但实际更多的瓶颈出在别的地方。
首先,源库的 binlog 写入是否有大事务?如果有大批量更新操作,binlog 会产生巨型事件,解析和序列化耗时会被拉长,这个属于源头瓶颈,只能通过拆分大事务解决。其次,Kafka 的 Topic 分区数是否足够?分区数太少,Consumer 的并发度上不去,积压就会出现。再者,Delta Lake 写入的 batch 配置是否合理?你设置的 15 分钟合并一次文件,那延迟就至少是 15 分钟起步。
排查的路径先看数据源 → 再看消息队列 → 最后看写入端,别上来就调 Flink 的并行度,不然很容易做了无用功。我自己常用的工具是 Kafka 的消费延迟监控(通过 Consumer Group 的 lag 指标),加上 Flink 的 Checkpoint 耗时监控,这两个数据基本能定位到 80% 的延迟问题。
5. 我在实际项目里总结的那点体会
5.1 增量层到底该做得复杂还是简单
关于 Delta Layer 的定位,不同团队有不同的做法。有的团队把增量层做得很重,在写入时做了字段清洗、标准格式化、甚至还做了预聚合。我个人的倾向是,Delta Layer 做"轻"一点比较好,宁可把清洗逻辑放到下游 DWD,Delta Layer 就老老实实保证数据的完整性和准确度。
原因也简单。增量层承接的是实时变更事件,如果在这里做太多业务逻辑,调试难度会增加很多。流式任务一旦逻辑复杂,数据水印、状态管理、乱序处理的坑都会接踵而来。保持轻量,既能快速发现问题所在,也让链路更加清晰。
5.2 监控比搭建更重要
我把增量层搭好之后,做的第一件事就是把监控补全。每日必须看的三个监控指标是:同步任务延迟、脏数据数量、表数据与源库主键对比差异。一旦最后一个指标出现持续差异,基本就意味着同步出现了数据丢失或重复,需要尽快干预。
补一个具体的监控配置思路。在 Flink 任务里,用MetricsReporter把currentEmitEventTimeLag、numRecordsIn、numRecordsOut上报到 Prometheus,再配置 AlertManager 做延迟告警。同时,写一个每日的数据对账 Spark 作业,对比源库和 Delta Table 主键集合的差集,生成差异报告。有了这两层,增量层基本是"半自动驾驶"的状态,不依赖人工巡检。
5.3 最后分享一个小技巧
前面写了很多方案和参数,最后分享一个小技巧,算是我这几年折腾 Delta Layer 的一点私货。
在做增量同步的表结构时,我强烈建议在每个表里都加上op_ts和event_type这两个字段,哪怕你现在觉得用不上。op_ts是事件发生的时间戳,event_type是事件类型。等你在哪个凌晨被电话叫醒,说"某张表的某一条数据被改坏了需要排查"的时候,这两个字段会成为你最有力的侦察工具——直接过滤出这一条数据的所有历史变更,看一眼事件时间线,基本就知道是怎么回事了。不要等到出了问题才后悔当初没加这两个字段,这是一个成本极低、收益极高的小设计。
Delta Layer 这个体系本身并不复杂,它的难点在于如何把增量捕获、消息队列、流式计算、数据落地这几段链路整合成一个稳定可靠的生产系统。把这套逻辑吃透了,实时数仓的基础也就扎实了一大半。希望这篇文章能帮少踩一些坑,如果你在实践过程中有其他问题,欢迎按照实际的链路和日志去逐一排查,大多数问题都是能被定位到具体环节的。