ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

实时数仓到底怎么建?从Kafka采集到指标计算、BI展示全链路拆解

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