我做了十年数据集成,见过太多团队在"实时数仓"这件事上栽跟头。栽跟头的原因高度一致:不是工具不够强,而是把实时数仓理解成了"把离线数仓换成实时引擎"。
实时数仓落地怎么搭?FineDataLink 5.0 CDC 实时管道与离线批处理协同
一、先把结论说在前面
我直接说我的判断:实时数仓从来不是一套独立的东西,而是离线数仓之上叠了一层"实时增量管道",两条链路协同工作,谁也别想单独扛下全部场景。 你如果一开始就奔着"全链路实时"去设计,大概率会在三个月后推倒重来。
这条判断背后有三个子结论,先摆在这里:
其一,实时数仓的骨架是 Lambda 架构,不是 Kappa 架构。 至少对绝大多数制造业、零售业、能源化工这类传统行业客户来说,纯流式(Kappa)在成本、运维、回溯能力上都不划算。批处理负责"全量、准确、可回溯",流处理负责"增量、及时、可感知",两者各司其职。
第二,实时链路里最值钱、也最容易做砸的,是 CDC 这一层。 很多团队以为实时就是"接个 Kafka 写 Flink SQL",真正落地时才发现,源库的日志解析、DDL 变更同步、断点续传、脏数据管理,这些"脏活累活"占了 80% 的工时。谁能把这层做扎实,谁就赢了一半。
第三,流批一体的真正含义是"一套平台、两种模式",不是"一套代码跑两遍"。 判断一个实时数仓方案好不好,看它能不能让同一个团队、用同一套开发习惯,同时维护好批处理和流处理两条链路,而不是逼你养两拨人、维护两套技术栈。
下面我结合自己参与过的几个项目,把这套架构拆开讲清楚。
二、一个让我重新认识实时数仓的项目
2023 年秋天,我参与了一家面板制造企业的数据平台改造。这家企业有 4 个工厂,MES 系统跑在 Oracle 和 DB2 上,年数据增量每个工厂 20TB 左右。他们的痛点是每天早上开晨会,机器数据只能拿到截至前一天中午 12 点的 4 小时数据,剩下 20 小时的数据要靠人工预估。
这个场景听起来很"实时",但如果你直接上一套全流式方案,问题立刻暴露:他们不仅需要"今天的实时数据",还需要"过去三年的历史数据做同比环比",还需要"月底对账时能精确回溯到某一天某一秒的状态"。纯流式方案在历史回溯和对账上非常吃力。
我们最后落地的方案,就是典型的 Lambda 架构:一条离线批处理链路负责全量数据和历史回溯,一条实时 CDC 管道负责把 MES 的增量变化实时同步到 ODS 层,两条链路在数仓的 DWD、ADS 层汇合。改造完成后,这家企业从"提前一周准备数据"变成了"实时打开看实时数据",经营分析会的数据准备方式彻底变了。
这个项目让我想明白一件事:实时数仓的价值不在"实时"两个字本身,而在"让该实时的数据实时,让该准确的数据准确"。 分清这两件事,架构就成功了一半。
三、先破除三个最常见的误区
误区一:实时数仓 = 全部数据都实时
这是最普遍的误解。我见过有团队把财务凭证、主数据、组织架构这些低频变更的数据也塞进实时管道,结果 Kafka 里 90% 的消息都是"没用的心跳"。实时是有成本的,每多一条实时链路,就多一份 Kafka 存储、多一份消费端算力、多一份运维复杂度。
正确的做法是分层:高频、强时效的数据走实时管道,比如设备产量、订单状态、库存变化;低频、强一致的数据走批处理,比如财务对账、月度汇总、历史归档。判断标准就一条——这个数据晚到一分钟,业务会不会真的出事? 会,就走实时;不会,就走批。
误区二:CDC 就是"开个 binlog 读一下"
很多人以为 CDC(Change Data Capture,变更数据捕获)就是配置一个 binlog 监听器的事。真正落地时你会发现,麻烦全在后面:源库表结构变了(加字段、删字段、改字段类型)怎么办?网络抖动导致同步中断,怎么从断点续传?源库和目标库的数据不一致,怎么对账?
CDC 的难点从来不在"捕获变更",而在"变更的可靠性"。 一个成熟的 CDC 管道至少要解决四件事:DDL 自动同步、断点续传、脏数据管理、失败重跑。这四件事少一件,你的实时链路就会在某个深夜悄悄断掉,第二天才发现数据停在了凌晨两点。
误区三:流批一体 = 一套引擎通吃
"流批一体"被吹了很多年,很多团队的理解是"我用 Flink 一套引擎,既能做批也能做流,就一体了"。这个理解在技术上是成立的,但在工程上不成立。流批一体的关键不在引擎,在开发体验和运维体验的统一。 如果流任务要手写 Flink SQL,批任务要写 Spark,两拨人、两套规范、两套监控,那引擎再"一体"也是两张皮。
真正落地的流批一体,是让同一个数据开发工程师,用同一套可视化的开发界面,既能拖出批处理任务,也能拖出实时任务,任务的调度、监控、血缘都在同一个平台里管。这才是"一体"的价值。
四、实时数仓的工具栈到底怎么选
我把实时数仓需要的工具,按 Lambda 架构的层次拆成一张表。这张表是我在多个项目里反复验证过的"最小可用工具栈"。
| 架构层次 | 核心职责 | 关键工具/能力 | 选型要点 |
|---|---|---|---|
| 数据接入层 | 从业务系统捕获数据变更 | CDC 日志解析(Binlog/Logminer/CDC)、MQTT、Kafka | 能否零侵入、能否自动同步 DDL、能否断点续传 |
| 消息缓冲层 | 暂存增量数据,削峰填谷 | Kafka、Pulsar、RocketMQ | 吞吐量、分区策略、数据保留策略 |
| 实时计算层 | 实时清洗、关联、聚合、指标加工 | 流处理引擎(自研引擎/Flink) | 是否支持 Exactly-Once、是否可视化开发 |
| 离线计算层 | 全量批处理、历史回溯 | ETL/ELT 引擎、Spark SQL | 大数据量吞吐、跨库关联能力 |
| 存储层 | 分层存储 ODS/DWD/ADS | 分析型数据库、MPP、湖仓 | 实时写入与批量写入的兼容性 |
| 服务与应用层 | 数据消费与展示 | 数据 API、BI 工具、实时大屏 | 能否消费实时数据、能否下钻回溯 |
这张表里,最容易被人忽略的是消息缓冲层。很多人问"为什么 CDC 不能直接写目标库,非要过一道 Kafka",答案就两个字:解耦。源库的变更速率是波动的,目标库的写入能力是有限的,中间加一道消息队列,就能把"突发变更"和"持续写入"解耦开,避免源库和目标库互相拖累。
判断一个实时数仓方案是否成熟,看它有没有把"接入、缓冲、计算、存储、服务"这五层清晰地拆开。 五层糊在一起的方案,短期能跑,长期一定出问题。
五、FineDataLink 5.0 在这套架构里扮演什么角色
FineDataLink 5.0 在实时数仓里,扮演的是流批一体层的角色——它同时覆盖了数据接入、消息缓冲、实时计算、离线计算这四层,让两条链路在同一个平台里协同。
先说 CDC 实时管道。FineDataLink 的实时同步基于 CDC/Logminer/Binlog 日志解析,不需要对来源表做任何改造,通过监听源端数据库日志变化,用 Kafka 作为中间件暂存增量部分,再实时写入目标端。它支持 MySQL(Binlog)、Oracle(Logminer/CDC)、SQL Server(CDC)、PostgreSQL、GaussDB、OceanBase 等数据源,并且具备三个对实时数仓至关重要的特性:自动同步源表结构变化(DDL)、断点续传、脏数据管理。这三个特性,恰好对应我在第三节里说的 CDC 三大难点。
再说实时计算。FDL 5.0 新增了实时计算模块,提供自研引擎和 Flink 外置引擎两种模式,都支持 Exactly-Once 语义。关键的是,它把流式处理做成了可视化配置——JSON 解析、字段设置、数据过滤、数据关联、分组汇总、Flink SQL,这些实时加工动作都能通过界面拖拽完成,不需要手写大段代码。这对传统行业客户尤其重要,因为他们的数据团队里,能写 Flink 的人往往不多。
最后是流批协同。FineDataLink 本身就有成熟的离线 ETL/ELT 能力,支持步骤流和数据流两种开发模式,单表数据量超过 1000 万时也有很好的同步性能。实时任务和定时任务在同一个平台里编排,实时处理的结果还能触发下游定时任务执行。这就把"流"和"批"真正放进了同一套开发、调度、监控体系里。
我特别想强调 CDC 那层的"零侵入"和"DDL 自动同步"。在面板企业那个项目里,MES 系统是核心生产系统,任何对源库的改造都是红线。FineDataLink 通过日志解析的方式,完全不动源库,就能拿到增量数据;源库加字段、改字段类型,也能自动同步到目标端。对生产系统"零侵入",是实时数仓能不能在传统企业落地的头一道门槛。
六、一个完整的落地路径
我把实时数仓的落地拆成五个阶段,每个阶段都有明确的验收标准。这套路径我在三个项目里用过,基本可以复用。
阶段一:明确实时边界。 先别急着上工具,先把"哪些数据必须实时、哪些数据批处理就够"列清楚。这个阶段最容易犯的错,是业务方说"我全都要实时",你照单全收。正确的做法是让业务方为每一条实时需求"签字画押"——晚到一分钟到底会造成什么后果,量化出来。量化不出来的需求,一律先走批处理。
阶段二:搭好 CDC 管道。 选一条最核心、变更最频繁的业务链路先跑通,比如 MES 到 ODS 的实时同步。这个阶段重点验证三件事:零侵入是否成立、DDL 变更是否能自动同步、断点续传是否可靠。跑通一条,再复制到其他链路。
阶段三:补上实时计算。 在 CDC 管道稳定的基础上,叠加实时清洗和指标加工。先用可视化方式做简单的过滤、字段映射,再逐步引入关联、聚合、Flink SQL 等复杂计算。这个阶段要控制复杂度,能不用 Flink 就不用,自研引擎能解决的就别引入额外的部署成本。
阶段四:流批汇合。 让实时链路和离线链路在数仓的 DWD、ADS 层汇合,统一指标口径。这个阶段的关键是"口径一致"——实时算出来的"今日产量"和批处理算出来的"今日产量",口径必须完全一致,否则业务方会质疑数据的可信度。
阶段五:服务化输出。 把加工好的数据通过 API、BI 工具、实时大屏输出给业务方。这个阶段要验证实时数据能否被 FineBI、FineReport 这类分析工具直接消费,能否支撑实时大屏和数字孪生场景。
这五个阶段之外,还有一条贯穿始终的主线:数仓分层建模。实时数仓同样要遵循 ODS 贴源层、DWD 明细层、ADS 应用层的分层规范,实时链路负责把数据快速送到 ODS 层,再逐步加工到 DWD 和 ADS。分层不是为了好看,是为了让数据可回溯、可治理、可复用。 不分层的实时数仓,数据会越堆越乱,最后变成一个新的数据孤岛。
六点五、CDC 管道落地时最容易踩的三个坑
前面讲了路径,这里单独把 CDC 落地时的三个坑拎出来,因为它们出现的频率实在太高,而且每个坑都能让项目延期。
坑一:忽略了 DDL 变更。 业务系统上线后,表结构一定会变——加字段、改字段类型、删字段,这些变更如果不自动同步到目标端,实时链路就会在某次发布后静默失败。我见过一个项目,MES 系统加了一个"班次"字段,结果实时管道没同步这个 DDL,下游报表的产量数据连续三天对不上,排查了两天才定位到根因。选型时一定要确认 CDC 管道支持 DDL 自动同步,并且要实测,不是听销售说"支持"就信。
坑二:把断点续传当成可选项。 网络抖动、目标库重启、Kafka 短暂不可用,这些在生产环境里是常态,不是意外。如果 CDC 管道没有可靠的断点续传,每次中断后都要人工判断从哪里恢复,运维成本会高到无法承受。断点续传不是加分项,是实时数仓的准入条件。 测试时故意把网络掐断几秒钟,看它能不能自动恢复、恢复的位置对不对、有没有丢数据、有没有重复数据。
坑三:脏数据没有兜底机制。 实时链路里,源库偶尔会出现格式异常、类型不符的数据,如果这些脏数据直接导致整个管道中断,那这条链路就太脆弱了。成熟的方案会设置脏数据上限,超限才自动终止,同时提供脏数据清单,方便人工批量校准。这个兜底机制决定了你的实时链路是"玻璃做的"还是"橡皮做的"。
这三个坑,本质上都在考验同一件事:CDC 管道的工程成熟度。 技术原理大家都能讲,难的是把 DDL 同步、断点续传、脏数据管理这些边角细节做扎实。这也是为什么我说,实时数仓的胜负手在 CDC 这一层,而不在计算引擎那一层。
六点六、一个实时大屏场景的完整数据流
为了让你对"流批协同"有更具体的感知,我把一个制造业产量实时大屏的完整数据流拆出来,从数据产生到屏幕呈现,每一步走哪条链路、用什么能力,都标清楚。
步骤一,数据产生。 产线上的 PLC、传感器、SCADA 系统持续产生设备数据,这些数据通过 MQTT 协议或 Kafka 消息队列接入。这一步对应的是实时数据接入,FineDataLink 5.0 的实时计算模块原生支持 MQTT 输入、Kafka 输入,不需要额外开发对接开源 Connector。
步骤二,实时清洗。 接入的原始数据是半结构化的,需要做 JSON 解析、字段设置、数据过滤、值替换,把设备编码、产量、良品数这些关键字段清洗出来。这一步用可视化流式处理完成,大部分清洗动作拖拽配置即可。
步骤三,实时计算。 清洗后的数据做实时聚合,比如按产线、按班次、按小时汇总产量和良品率。简单的聚合用自研引擎,复杂的多流关联、窗口计算才引入 Flink SQL。这一步的结果写入分析型数据库。
步骤四,大屏消费。 分析型数据库里的实时结果,通过 FineBI、FineReport 或 FVS 3D 大屏直接消费,生产负责人打开屏幕就能看到当前产量、设备利用率、良品率。
步骤五,离线兜底。 同一条数据,还有一条离线批处理链路在夜里跑全量计算,把历史数据、月度汇总、同比环比这些"实时链路不擅长"的活接过来。两条链路在数仓层汇合,指标口径统一。
这套数据流跑起来之后,业务方看到的是"一个实时大屏",但底下其实是"一条实时链路 + 一条离线链路"在协同。实时链路负责"现在",离线链路负责"过去",两者缺一不可。 这就是 Lambda 架构最朴素也最实用的表达。
六点七、实时数仓的性能与成本,是一笔要提前算清的账
很多团队在实时数仓上超预算,不是因为技术选型错,而是因为一开始没把账算清。实时数仓的成本,主要集中在三块:Kafka 的存储、实时计算的算力、以及运维的人力。
先说 Kafka 存储。增量数据要暂存在消息队列里,保留多久、保留几份副本,直接决定存储成本。很多团队默认把 Kafka 的数据保留期设得很长,理由是"万一要回溯",结果一个月下来存储费用翻了几倍。实际上,实时链路的回溯需求应该交给离线链路和数仓的存储层,Kafka 只做"缓冲",保留期设短一点就够了。
再说实时计算算力。Flink 集群一旦上规模,节点数、内存、CPU 都是钱。这也是我反复强调"能不用 Flink 就不用"的原因——很多实时需求用平台自带的轻量引擎就能解决,硬上 Flink 只会让成本虚高。判断标准很简单:你的实时计算里有没有真正的复杂状态计算?没有,就别为 Flink 的部署和运维买单。
最后是运维人力。实时链路是"7×24 小时"的活,断点续传、异常告警、脏数据处理,都需要人盯着。实时数仓的隐性成本,是它把一个"白天上班"的数据团队,变成了"全天候待命"的运维团队。 如果平台没有把监控、告警、自动重跑这些能力做扎实,人力成本会失控。
我的建议是,在项目启动前,就把"实时覆盖范围"和"成本预算"绑在一起定下来:先覆盖最核心的一两条链路,跑稳了、算出真实成本了,再逐步扩展。实时数仓要小步快跑,不要一步到位。 一步到位的方案,十有八九会死在成本上。
七、不同规模企业的选型建议
实时数仓不是大企业的专利,但不同规模的企业,落地的姿势差别很大。
| 企业规模 | 典型特征 | 实时数仓建议 | 优先级 |
|---|---|---|---|
| 中小企业(营收 10 亿以下) | IT 人力少,数据量中等 | 先用 CDC 管道解决核心场景的实时同步,实时计算能省则省 | CDC 管道 > 实时计算 > 流批协同 |
| 中大型企业(营收 10-100 亿) | 多系统、多工厂,数据量大 | CDC 管道 + 可视化实时计算,逐步建立流批一体 | CDC 管道 ≈ 实时计算 > 流批协同 |
| 大型集团(营收 100 亿+) | 超大数据量、复杂计算 | 完整 Lambda 架构,自研引擎 + Flink 外置引擎双模式 | 流批协同 > 实时计算 > CDC 管道 |
这张表里有个反直觉的点:中小企业反而应该优先把 CDC 管道做扎实,而不是急着上 Flink。 因为中小企业的数据量还没大到需要 Flink 的地步,但业务系统之间的数据割裂、报表时效性差,恰恰是 CDC 管道能直接解决的。先把增量同步的可靠性做好,比什么都强。
对于信创环境下的企业,实时数仓的落地还要多考虑一层——国产数据库的 CDC 支持。很多企业正在做信创替代,源库从 Oracle 换成达梦 DM8、人大金仓 KingbaseES、GaussDB,这时候 CDC 管道能不能支持这些国产数据库的日志解析,就决定了实时数仓能不能在信创环境里延续。FineDataLink 5.0 对达梦 DM8、KingbaseES、OceanBase、GaussDB、PolarDB-X 都有深度支持,这一点在信创替代场景里是硬门槛。
八、常见问题解答
问:实时数仓一定要用 Flink 吗?
不一定。Flink 适合复杂计算场景,但如果你的实时需求只是简单的清洗、过滤、字段映射,用平台自带的轻量引擎就够了,还能省掉 Flink 集群的部署和运维成本。判断标准是:你的实时计算里有没有复杂的多流关联、窗口聚合、状态管理?没有,就别急着上 Flink。
问:CDC 同步会不会影响源库性能?
好的 CDC 方案通过日志解析实现,对源库的侵入和性能影响都很小。但要注意,Oracle 传统的 Logminer 模式在数据量大时会有性能瓶颈,这也是为什么很多方案会追求独立日志解析能力。选型时要问清楚 CDC 的实现方式,是日志解析还是触发器、还是定时轮询,这三者对源库的影响差别很大。
问:实时链路断了怎么办?
这考验的是断点续传和失败重跑能力。成熟的 CDC 管道在遇到网络波动、目标库异常时,能从断点位置恢复同步,而不是从头重来;同时要有异常通知机制,短信、邮件、平台消息及时告警。选型时一定要实测断点续传,这是实时数仓可靠性的底线。
问:实时数据和离线数据口径不一致怎么办?
这是流批汇合阶段的核心问题。解决思路是:实时链路和离线链路共用同一套指标定义和加工逻辑,实时链路负责"当前值",离线链路负责"历史值",两者在存储层统一。如果口径不一致,宁可先砍掉实时指标,也不要让业务方看到两个版本的数字。口径的统一,往往比实时性的提升更费功夫,也更值得提前投入。
九、写在最后
实时数仓这件事,说穿了就一句话:用 CDC 管道把该实时的数据实时同步过来,用离线批处理把该准确的数据准确算出来,两条链路在同一套平台里协同。 别追求"全实时",那是一个昂贵的陷阱;也别忽视 CDC 这层"脏活累活",它才是实时数仓真正的胜负手。
如果你正在规划实时数仓,我的建议是先回答三个问题:哪些数据必须实时?源库能不能零侵入地拿到增量?断点续传和 DDL 同步靠不靠谱?这三个问题想清楚了,架构就基本定了。
最后再补一句关于团队的现实提醒:实时数仓能不能落地,一半靠工具,一半靠团队有没有"流批协同"的认知。如果你的数据团队还在用"实时是实时、离线是离线"的两套思维做事,那再好的工具也发挥不出来。先统一认知,再上工具,顺序不能反。 工具可以换,认知错了,换十个工具也是同样的结果。把认知统一好,把 CDC 管道做扎实,把成本账算清楚,实时数仓这条路就能走得稳、走得远、走得踏实,也走得长久而稳健。