很多企业一提实时数仓,第一反应往往是:
Kafka怎么搭?
Flink要不要上?
延迟能不能做到秒级?
但真正做过项目以后会发现,技术组件反而不是最难的。
真正难的是:
数据从哪里来、怎样持续进来、进入以后怎么算、指标怎么统一、链路断了怎么恢复,以及最后业务看到的数字到底可不可信。
尤其是实时场景。
离线数仓算错了,可能第二天发现;
实时数仓算错了,错误可能几分钟内就已经出现在经营大屏、库存预警,甚至进入业务决策。
所以实时数仓真正要建设的,不只是一个“跑得更快的数仓”,而是一条完整的数据链路:
业务系统 → CDC / Kafka → 实时加工 → 数仓分层 → 指标计算 → BI展示 → 监控治理。
在正式展开之前,我整理了一套《数据仓库建设解决方案》,里面覆盖数据集成、数仓建设、数据治理等常见项目场景。如果正在做实时数仓、数据平台或者企业数据集成,可以结合下面这条完整链路一起参考。
需要自取:https://s.fanruan.com/7igmg(复制到浏览器)
一、实时数仓第一步,不是上Kafka,而是先定义“什么值得实时”
实时数仓最容易犯的第一个错误,就是:
所有数据都要求实时。
实际上,不同业务对时效性的要求完全不同。
支付风控可能要求秒级;
库存预警可能要求1分钟以内;
门店销售看板3分钟刷新一次已经够用;
月度利润、客户价值分层,本身就没有必要实时计算。
所以项目开始之前,最好先给不同数据定义一个时效SLA。
比如:
支付状态:30秒以内;
库存变化:1分钟以内;
订单汇总:3分钟以内;
经营指标:5分钟以内;
财务结算:T+1。
这一步很重要。
因为它直接决定后面到底应该使用:
实时流、准实时任务,还是传统批处理。
成熟的数据平台往往不是“全实时”,而是:
实时链路 + 准实时链路 + 离线链路并存。
确定时效以后,才进入采集。
业务数据库里的订单、库存、合同,可以通过CDC持续捕获数据变化;
已经事件化的业务系统,可以直接把订单、支付、退款等事件写入Kafka;
IoT、设备、日志数据,则可能通过MQTT、消息队列等方式进入实时链路。
这一层有一个非常重要的原则:
实时采集的重点,不是不断查询“现在有什么”,而是持续捕获“刚刚发生了什么变化”。
比如一张订单表已经有5亿条数据。
如果每分钟都通过更新时间重新扫描,数据量越大,对源库的压力越明显。
而CDC关注的是:
订单10001从待支付变成已支付。
两种方式看起来都能拿到结果,但背后的系统成本完全不同。
真正到了企业环境,问题还会继续扩大:ERP一套库、CRM一套库、MES又是一套库,还有Kafka、API和文件。
如果每一个来源都单独写采集程序,后面维护的其实不是数仓,而是一堆接口。
像这类场景,更常见的做法是先把数据接入这一层统一起来。比如用FineDataLink 5.0连接业务数据库、Kafka等数据源,把不同系统里的增量变化持续送到后面的数仓链路里。后续无论是增加数据源、调整同步目标,还是查看链路状态、排查异常,都不需要再到几十套独立脚本里逐个处理。
对于实时数仓来说,这一步的意义很直接:
先把“数据怎么稳定进来”这件事标准化,后面的实时计算和数仓建模才有稳定的数据入口。
二、Kafka不是实时数仓,它解决的是“数据怎么流动”
很多实时数仓架构图里,Kafka都放在最中间。
于是很容易产生一种误解:
有Kafka,就有实时数仓。
实际上,Kafka更像整个系统里的:
实时数据总线。
假设订单系统每秒产生1000笔订单。
如果订单系统直接连接:
风控系统、营销系统、实时大屏、数仓、推荐系统……
很快就会形成大量点对点接口。
引入Kafka以后,可以变成:
业务系统 → Kafka → 多个消费者。
上游只负责生产事件,下游按照自己的需求消费。
它主要解决两个问题:
缓冲和解耦。
但Kafka真正落地时,还要提前设计几个问题。
Topic不能乱建
最好围绕:
业务域 + 事件类型
进行规划。
例如:
order_event
payment_event
refund_event
inventory_event
而不是所有业务都塞进一个“大Topic”。
Partition Key决定局部顺序
假设一张订单先支付,后退款。
如果两个事件随机进入不同分区,就可能出现消费顺序不一致。
所以订单类场景经常按照:
order_id
进行分区,让同一个业务对象尽可能保持局部有序。
Kafka最好保存业务事件,而不是报表结果
应该保存:
订单10001在10:23支付成功,金额299元。
而不是直接保存:
华东区销售额增加299元。
因为前者可以继续被风控、推荐、营销和数仓复用;
后者已经绑定在某一个具体分析口径上。
所以Kafka解决的是:
让数据稳定地流动起来。
而真正把数据变成分析结果,还要依赖后面的实时加工和数仓建模。
三、实时加工真正难的,不是SUM,而是时间、状态和一致性
很多实时计算Demo看起来都很简单:
SUM(order_amount) GROUP BY region
但生产环境远没有这么简单。
一笔订单可能经历:
创建 → 支付 → 修改 → 退款 → 取消退款。
如果每一次变化都直接把订单金额累加一次,结果一定会错。
所以实时加工面对的并不是一张静态表,而是:
持续变化的事件和状态。
重复问题
实时系统经常采用“至少一次”投递机制。
网络异常、任务重启、消费失败,都可能导致一条消息重新处理。
所以必须考虑:
业务唯一键和幂等。
也就是说:
同一个业务事件即使处理两次,最终结果也不能多算一次。
乱序和迟到
业务事件发生的顺序,并不一定等于系统收到的顺序。
10:01发生的支付事件,可能10:03才到;
10:02的订单修改反而先被消费。
所以实时计算必须区分:
Event Time,事件发生时间
和
Processing Time,系统处理时间。
尤其做分钟、小时级指标时,如果这个问题没有处理好,同一份数据的实时结果和离线结果很容易长期对不上。
维度关联
订单流里通常只有:
product_id
但经营分析真正需要的是:
商品、品牌、品类、事业部、区域。
于是实时事实流必须关联维度数据。
更麻烦的是,维度自己也会变化。
比如某商品1月份属于A品类,3月份调整成B品类。
那1月份历史订单到底应该展示A,还是跟着变成B?
这已经不是简单的JOIN问题,而是:
历史维度和当前维度的口径设计。
实时聚合
经营层最终不会看一条条消息,而是看:
今日销售额、实时订单量、区域GMV、商品销量、良品率。
所以数据还要不断经历:
过滤、转换、关联、聚合。
实际项目里,这里的加工逻辑也要分复杂度。
像字段转换、条件过滤、维表关联、分组汇总这类比较固定的规则,没有必要每一次都重新开发一套流计算程序。
例如订单数据进来以后,需要先补充商品、区域、组织信息,再按照区域实时汇总销售额;设备数据进来以后,要先过滤异常值,再按产线统计产量和良品率。这类链路可以直接放在FineDataLink 5.0里完成必要的数据转换、关联和聚合,再把加工后的结果继续送往下游。
而真正涉及复杂状态、窗口计算、乱序处理或者大规模流计算时,再交给专门的实时计算引擎。
这样整个实时加工层不会变成“什么都上Flink”,而是根据计算复杂度拆开:
标准的数据处理走固定链路,复杂实时计算再单独设计。
四、实时数仓照样要分层,不要把Kafka直接接到所有报表
有些项目为了追求快,会设计成:
Kafka → 指标表 → BI。
刚开始业务少的时候,确实很快。
但业务一多,问题马上出现。
销售大屏计算一次销售额;
运营看板又计算一次;
财务分析再计算一次。
最后同一个指标在三套任务里存在三套逻辑。
所以实时数仓依然需要分层。
比较常见的逻辑是:
ODS → DWD → DWS → ADS
ODS:保留原始变化
ODS的核心作用可以概括成三个字:
留现场。
订单创建、修改、支付、退款这些原始变化尽量完整保存下来。
未来出了问题,才能判断:
源数据错了,还是加工错了。
如果原始事件都没有保留,后面的排错只能靠猜。
DWD:统一业务事实
DWD开始真正做标准化。
包括:
主键统一、字段统一、编码统一、状态统一、业务含义统一。
例如多个渠道都有订单,可以最终统一成:
订单事实表、支付事实表、退款事实表。
这一层真正解决的是:
下游到底应该相信哪一份业务数据。
DWS:沉淀公共能力
比如:
区域小时销售额、
客户累计消费额、
商品日销量、
设备小时产量。
这些数据如果会被多个应用反复使用,就应该提前沉淀。
否则每个应用都从DWD开始重新计算,公共逻辑就会不断复制。
ADS:面向具体业务场景
ADS才真正服务:
经营驾驶舱、实时大屏、风险预警、专题分析。
所以一个比较实用的原则是:
公共逻辑往下沉,个性逻辑往上放。
项目做到这里以后,很容易出现另一个问题:
同一份数据,被不同下游重复读取、重复加工。
例如订单明细进入DWD以后,经营分析要按区域汇总,销售分析要按商品汇总,库存系统还要继续使用部分订单信息。
如果每个场景都重新从源库拉一遍,数据链路很快会越来越多。
这时候可以利用FineDataLink 5.0的数据分发思路,把同一条上游数据先做统一处理,再根据不同用途送往不同目标:明细数据进入DWD,公共汇总进入DWS,需要继续消费的数据再发送给其他系统。
这样做比“每个应用各自采一遍数据”更重要的一点是:
大家尽量从同一份标准数据继续往下加工。
一旦上游字段、状态或者业务规则发生变化,也不至于在十几条独立链路里分别修改,实时数仓的分层才真正有复用价值。
五、指标层必须提前统一,不能把数据都扔给BI现场算
实时数仓最后真正交付给业务的,其实不是数据库表。
而是:
指标。
比如管理层看到:
今日GMV:1280万元。
这个数字背后其实藏着一整套业务规则:
哪些订单状态算成交?
未支付订单算不算?
部分退款怎么处理?
优惠券是否计入GMV?
跨天退款算哪一天?
测试订单是否剔除?
如果这些问题不提前统一,数据即使10秒刷新一次,也没有意义。
所以更合理的链路应该是:
原始事件 → 标准事实 → 公共汇总 → 指标 → BI。
例如:
净支付金额 = 有效支付金额 - 有效退款金额
先把“有效支付”和“有效退款”的定义统一。
然后再根据:
日期、区域、渠道、产品
生成各种派生指标。
指标体系还可以继续分为:
原子指标、派生指标、复合指标。
例如:
支付金额,是原子指标;
华东区今日支付金额,是限定了区域和时间的派生指标;
客单价 = 支付金额 ÷ 支付客户数,是复合指标。
这样做的意义在于:
指标的定义位置前移。
而不是等数据到了BI,分析人员再临时解释。
到了这一层,重点已经不是“数据能不能过来”,而是:
BI最终拿到的数据,能不能直接用于分析。
如果底层只是把订单、退款、客户、商品几张明细表原样同步过去,那么大量清洗、关联和口径判断还是会堆到报表端。
更合理的做法,是在前面的数据链路里就把可以提前确定的逻辑处理掉。
比如通过FineDataLink 5.0把订单、退款、商品等数据持续汇入分析库之前,先完成必要的字段整理、数据关联和公共汇总,再让BI读取已经整理好的结果表。
这样经营驾驶舱看到“今日销售额”“区域完成率”时,不需要每张报表重新解释一次底层业务数据。
整个链路也会更清楚:
先把分散数据接进来并完成必要加工,数仓统一模型和指标口径,BI再做筛选、下钻、对比和展示。
最终业务看到的是分析结果,而不是底层数据加工过程。
六、实时数仓上线以后,至少盯住这5个问题
很多实时系统最危险的情况并不是:
任务挂了。
而是:
任务看起来还在运行,但数据已经错了。
所以真正进入生产环境以后,至少要持续监控五件事。
端到端延迟
不能只看Kafka有没有积压。
真正应该监控的是:
业务系统发生变化 → BI最终可见
整个链路用了多久。
因为消费者延迟5秒没有意义,如果后面的聚合、写入和BI刷新又用了20分钟。
消费Lag
当生产速度大于消费速度时,Kafka会开始积压。
系统可能没有任何报错,但所谓实时数据已经慢慢变成:
10分钟前的数据。
所以Lag本质上监控的是:
链路有没有越来越追不上业务变化。
数据量是否闭合
最好建立一条完整对账关系:
源端变化量 → 消息量 → 消费量 → 加工量 → 写入量 → 异常量
比如源系统新增100万条变化记录,最终只有98万条进入数仓。
那另外2万条必须能够解释。
否则系统表面上一直在运行,实际上可能早就开始丢数据。
重复和异常
实时链路经常遇到:
重复消息、失败重试、断点恢复、脏数据。
所以要持续监控:
业务主键重复率、异常记录数量、失败重试次数。
不能只看“任务成功”。
因为一个任务显示成功,并不代表最终数字一定正确。
实时和离线结果是否一致
很多企业都会同时保留:
实时结果 + T+1离线结果。
这其实是一种很好的校准机制。
比如实时GMV显示1000万,第二天离线重新计算只有950万。
不能简单解释成:
“实时有误差很正常。”
要继续拆:
到底是迟到数据?
重复事件?
退款状态变化?
还是指标口径本身不一致?
真正成熟的实时数仓,并不是永远不出问题。
而是:
出了问题以后,可以快速知道问题发生在哪一层、影响了多少数据,以及应该从哪里重新恢复。
结语
实时数仓看起来是很多技术组件的组合:
CDC、Kafka、Flink、数据仓库、指标体系、BI。
但如果把这些技术名词全部去掉,它真正要解决的问题其实非常简单:
把业务世界持续发生的变化,及时转化成可信、统一、可以直接使用的数据。
完整链路可以概括为:
业务产生变化 → CDC / Kafka采集 → 清洗、去重、关联、聚合 → ODS / DWD / DWS / ADS分层 → 指标统一计算 → BI分析展示 → 延迟、Lag、对账与异常监控。
所以企业建设实时数仓时,真正值得关注的并不是:
延迟能不能做到1秒。
而是:
哪些数据真的值得实时?
整个链路出现异常以后能不能追溯和恢复?
同一个指标在不同系统、不同看板里,能不能始终保持同一个业务含义?
只有这三个问题解决以后,Kafka才不只是消息队列,实时计算才不只是技术能力,BI也不只是最后那块大屏。
它们才真正组成一套能够服务经营决策的实时数据体系。