ARTICLE DETAIL

资讯详情

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

电商实时处理架构全解析:从Flink到Kafka的链路设计

电商实时处理架构全解析:从Flink到Kafka的链路设计 去年大促我蹲在会议室盯订单实时看板屏幕上GMV数字卡在一个值将近十分钟没动。运营同事跑过来问是不是链路挂了我查了一圈才发现上游订单库的Binlog消费lag飙到了几十万条Kafka topic的分区消费者被一个慢算子拖住了整个实时链路从上到下都在等。那次之后我彻底明白电商数据实时处理架构不是装个Flink就能交差的它是一整条从数据产生到最终可见的流水线任何一环出问题大屏上那个数字就是骗人的。这篇文章想把电商实时处理架构里里外外拆一遍。无论你是刚接触实时计算的数据工程师还是已经在用Flink但想把链路设计得更稳的老手都可以从里面找到能直接落地的思路消息队列怎么选、窗口参数怎么设、实时数仓怎么分层、线上遇到背压和数据倾斜怎么办。最后我会给出一条订单实时统计的完整链路案例你可以照着思路搭一套能跑通的最小实现。1. 实时处理解决的是电商里哪些具体问题1.1 离线数仓最大的短板是时间差做电商数据的人应该都有过这种经历运营早上要昨天一整天的销售明细离线任务凌晨两点跑完九点才能看到数据。T1的节奏对财务对账、月度复盘来说没什么问题但遇到大促或者突发流量等数据跑完再做决策就晚了。我见过最典型的场景是活动商品的价格调整。运营希望根据实时的转化率调整优惠券发放力度如果只能看昨天的数据那今天的流量就浪费了。实时处理解决的不是计算速度问题而是决策时效问题。同样一份GMV离线算和实时算数学结果应该是一样但用途完全不一样。实时的价值在于它把事后复盘变成了事中干预这是电商业务最核心的诉求。还有一个容易被忽略的点离线数仓处理的是已经结束的数据天然没有延迟概念。但实时处理面对的是无限流数据一直来窗口什么时候关、乱序数据怎么处理、重复数据怎么去重都是离线根本不会遇到的问题。这套复杂度才是实时架构真正要去解决的。1.2 电商里典型的实时场景盘点电商的实时场景比很多人想象的多而且每个场景对延迟的要求完全不同大促大屏和经营看板分钟级延迟就够用了主要看GMV、订单量、客单价的走势。这类场景关心的不是单笔数据而是聚合趋势用流处理引擎做窗口聚合最合适。实时库存扣减秒级甚至要求准实时。库存超卖是电商大忌用户下单之后发现没货体验极差。库存链路通常不走纯流计算而是依赖Redis或者数据库事务但这部分数据也要回流到实时数仓里用于分析。交易风控毫秒级比如支付的时候判断设备、IP、行为序列是否异常。风控链路往往是独立系统实时架构要做的更多是把风控结果和订单数据做关联方便后续分析。推荐和搜索的实时特征用户刚点击了某个商品下次请求就想看到相关内容。这个链路的数据从行为日志来经过实时计算生成特征再写入特征服务延迟要求通常在秒级以内。实时对账和异常监控支付流水、退款流水、结算数据之间可能存在不一致实时对账能立刻发现差异。这个场景不需要很高的延迟但对准确性极其敏感一条数据错了都要捞出来。这些场景背后共享同一套底层设施日志采集、消息队列、流计算引擎、实时存储。架构设计得好这些需求都能复用同一张实时数据管道而不是每个业务各搭一套。1.3 实时的粒度其实分三档我经常被问到一个问题实时到底要快到什么程度这个问题没有标准答案得按业务来分。毫秒级秒杀、风控、支付结果这类链路要求数据产生后几乎即刻被消费。这类场景通常不适合用Flink做全链路因为Flink的调度和窗口也会引入几十毫秒到秒级的开销。业界通常用直连服务、Redis、消息队列的即时消费来实现。秒级库存、价格、个性化推荐。这类场景对实时性要求高但允许轻微的延迟波动。Flink处理链路如果优化得好秒级是可以做到的关键在于避免背压和频繁checkpoint引起的停顿。分钟级经营看板、大促大屏、实时报表。这是最典型的准实时场景用Flink的窗口聚合或者微批处理就能覆盖成本和复杂度和毫秒级完全不是一个量级。很多团队一开始就把目标定在秒级甚至毫秒级结果把架构搞得极复杂最后发现运营看板根本不需要这个精度。我的经验是先明确业务对延迟的真实容忍度再用最小成本去满足。实时不是越快越好而是刚刚好。2. 一条实时数据从产生到可见中间要过四层2.1 采集层埋点、日志和数据库变更缺一不可数据源是实时架构的地基。电商实时数据的来源主要有三类第一类是客户端埋点数据。APP和H5页面里的浏览、点击、加购行为通过SDK上报到服务端再由服务端转发到消息队列。埋点的关键是全链路ID不管是设备ID、用户ID还是会话ID必须能从一次点击一路追溯到最终的订单否则后面的行为分析和转化漏斗根本做不准。第二类是服务端日志。比如订单服务、支付服务打印的业务日志包含订单创建、支付成功、退款发起等关键事件。这类日志通常已经有比较规范的格式关键是要把时间戳统一最好使用事件发生时间而不是日志落地时间否则乱序问题会在计算层放大。第三类是数据库变更日志也就是CDC。订单在MySQL里insert了一条记录库存表update了一个数量这些变更本身就是一个事件流。用Canal或者Debezium监听Binlog把变更转成结构化的消息发到Kafka实时计算就能感知数据库的变化。这条路径在实时数仓里越来越重要因为订单、支付这类核心交易数据最终的准确口径还是在数据库里。采集层最容易踩的坑是重复上报和丢失上报。客户端网络抖动导致重试会造成日志重复服务端宕机可能导致日志没来得及flush。这些问题不能指望在采集层完全解决但要在设计时就给计算层预留幂等处理的能力后面才好兜底。2.2 缓冲传输层Kafka为什么是标配采集层的数据不会直接进计算引擎中间必须有一层消息队列缓冲。Kafka在电商实时处理架构里基本是事实标准原因其实不复杂。首先是削峰填谷。大促零点的流量可能是平峰的几十倍计算引擎如果直接扛这个流量要么按峰值扩容导致平时资源浪费要么直接被冲垮。Kafka把数据先存下来计算引擎按自己能处理的速度消费中间就多了一个缓冲。其次是解耦。上游系统不需要知道下游有几个消费者。订单数据进了Kafka实时数仓消费一份风控系统消费一份推荐系统消费一份互不干扰。新增一个消费者不需要改动上游这是架构演进特别需要的灵活性。选型上主要考虑Kafka、Pulsar、RocketMQ这几类。Kafka的生态最成熟、吞吐量高对大数据量级的日志和CDC支持很好缺点是分区扩容和Rebalance的运维体验一般。Pulsar的存算分离架构在弹性扩容上有优势但整体运维复杂度比Kafka高。RocketMQ在事务消息上有自己的特色适合和业务强绑定的消息场景。我的建议是以日志和CDC为主、数据量大的链路用Kafka业务消息、需要事务保障的用RocketMQ或者Pulsar不要一套队列包打天下。Kafka使用时有一个关键参数要关注分区数。分区是并行度的上限Flink的并行度、消费者数量都受分区数约束。分区太少会限制吞吐太多会带来文件句柄和Rebalance开销。一般建议根据峰值吞吐量、单分区吞吐量以及下游并行度来综合估算而不是拍脑袋定个数。2.3 计算层流处理引擎负责干活数据从Kafka消费出来之后就到了实时计算层。这里做的活包括数据清洗过滤无效数据、补全字段、数据转换JSON解析、格式统一、数据关联订单流和商品维表关联、数据聚合按窗口计算GMV、订单量、以及把结果写到下游存储。目前主流选择是Apache Flink。Flink的优势在于真正的流式计算模型处理延迟低状态管理成熟checkpoint机制让计算可以从故障中恢复而不丢数据。Spark Structured Streaming在微批模式下也能实现秒级延迟如果团队Spark技术栈很强也可以作为备选但在需要低延迟和复杂事件处理的场景里Flink显然更顺手。选计算引擎的时候团队的技术积累比引擎的纸面性能更重要。我见过团队硬上Flink结果没人会调状态后端checkpoint频繁失败最后又退回Spark。工具没有绝对好坏关键是有没有人能把它用好。2.4 存储与查询层结果要能快速被读计算完成后结果数据必须落到一个面向查询的存储里。这一步很多人会忽略以为Flink算完就结束了。实际上计算结果要支撑大屏、报表、在线服务对查询能力和并发要求是完全不同的。Redis适合存单个维度的实时值比如当前订单总量、当前GMV大屏直接GET就能拿到。优点是快缺点是无法做多维分析。ClickHouse / Doris适合存明细和轻度聚合结果支持SQL查询能做多维OLAP分析。ClickHouse查询极快但不适合高频单点更新更适合批量写入后分析Doris在实时更新和主键模型上更友好最近几年在实时数仓里用得很广。Elasticsearch适合订单明细的搜索和过滤比如查某个用户近一小时的订单。它的聚合能力不如ClickHouse但全文检索和条件过滤很强。HBase适合按主键查询的超大明细表但在电商分析场景里用得越来越少因为写链路重、查询能力单一。链路设计上通常会同时写多个存储聚合结果写Redis供大屏读取明细和轻度聚合写ClickHouse/Doris供分析师查询。这样一套数据不同出口各取所需。层级核心组件主要职责延迟量级采集SDK、Logstash、Canal、Debezium埋点、日志、Binlog捕获百毫秒缓冲Kafka / Pulsar削峰、解耦、数据积压毫秒级计算Flink / Spark Structured Streaming清洗、关联、窗口聚合秒级存储/查询Redis、ClickHouse、Doris、ES大屏读取、OLAP分析毫秒~秒级每个新项目最有效的切入方式就是先把这张表的每一层选型定了再往下设计细节。层与层的边界越清晰出问题的时候越容易定位。3. Flink实时计算的三个核心概念弄懂它们你才算真正入门3.1 窗口把无限流切成有限块流数据是无限的但业务指标是按时间范围算的比如最近5分钟的下单量今日累计GMV。Flink里就引入了窗口的概念把无限数据流按时间切成一块块有限数据。常见的有三种窗口滚动窗口TUMBLE固定长度且不重叠适合算每分钟的订单量、每小时的销量滑动窗口HOP有固定的长度和滑动步长适合算最近5分钟最近1小时这种持续滚动的指标会话窗口SESSION根据不活跃间隔切分适合分析用户一次完整访问的行为序列。用Flink SQL开一个1分钟滚动窗口算订单量代码很直观INSERT INTO dws_order_1min SELECT order_time, COUNT(*) AS order_count, SUM(pay_amount) AS gmv FROM dwd_order_paid GROUP BY TUMBLE(order_time, INTERVAL 1 MINUTE);窗口参数一旦设错指标含义就错了。做电商看板的人应该深有体会页面显示实时销售和昨日至今天累计这两类指标用的窗口逻辑完全不同。实时销售可能用10秒滚动窗口再刷新展示累计则要把窗口内数据累到一个带状态的指标里。所以设计窗口时先问清楚业务到底要哪个口径再动手写SQL。3.2 水位线处理乱序数据的核心机制流处理里最让人头疼的问题是乱序。用户点击发生在前日志反而后到订单创建时间比支付时间晚到但逻辑上创建必须在先。如果严格按到达顺序计算指标必然出错。Flink的解决方案是事件时间加水位线。事件时间就是业务发生的时间日志里自带的时间戳水位线则表示当前事件时间推进到了某个点早于这个点的数据应该已经到了。窗口触发条件会基于水位线判断水位线越过了窗口结束时间窗口就触发计算。这里有个永恒的矛盾水位线设得越宽松容忍乱序能力越强但窗口计算等得越久结果越慢设得越紧结果越快但迟到的数据会被丢弃。实际项目里我一般会根据业务容忍度设置迟到数据允许的延迟时间比如订单流设置30秒到1分钟的水位线延迟再单独处理极端迟到的数据。Flink SQL里通过WITH参数指定事件时间和水位线比如CREATE TABLE dwd_order_paid ( order_id BIGINT, pay_amount DECIMAL(10,2), pay_time TIMESTAMP(3), WATERMARK FOR pay_time AS pay_time - INTERVAL 30 SECOND ) WITH (...);这个30秒就表示最多容忍30秒的乱序。超过这个范围的数据会被丢弃或者被发到侧输出流里另行处理。做实时架构必须把允许丢多少数据和允许慢多少秒当成两个明确的产品参数而不是靠工程师拍脑袋。3.3 状态、检查点与精确一次流计算里的很多操作依赖状态去重需要记住已经见过的ID计数需要保存当前的累计值窗口聚合需要缓存窗口内的数据。Flink把状态抽象成可以持久化的存储并且通过检查点机制定期把状态快照保存到外部系统。检查点做的是异步快照Flink会周期性触发全链路的状态快照。任务崩溃后可以从最近一次成功的检查点恢复相当于把流处理回滚到了某个时间点然后重新消费。这样一来即使后续程序崩溃也不会疯掉。在这个基础上Flink提供了精确一次Exactly-once的语义。注意这个一次性是针对Flink内部的要实现端到端的精确一次还需要Kafka的消费者、下游Sink都配合事务机制。Kafka可以配合实现跨系统的exactly-once但开销不小很多电商场景其实用的是至少一次加幂等的方案允许重复读取但写下游时通过主键或去重逻辑保证数据不重复消费。这里要提醒一个细节checkpoint开启后同一时间的恢复能力是有代价的状态越大checkpoint越频繁恢复越慢。实践中建议根据业务重要性设置checkpoint间隔比如核心交易链路30秒一次日志分析链路可以放宽到2分钟。全链路每30秒checkpoint一次的压力真不小状态超过10GB的作业要仔细调优。3.4 CDC链路数据库变更也能实时参与计算CDC变更数据捕获是实时数仓里极具分量的一环。订单、支付这类核心数据最终准确口径在数据库里但离线同步太慢直接查库又影响业务。CDC思路是让业务数据库的Binlog作为数据源实时把增删改的变化发到Kafka下游Flink拿到之后直接做计算。链路大概是MySQL Binlog - Canal/Debezium - Kafka - Flink。Canal是阿里开源的老牌方案对MySQL支持好部署简单Debezium支持多种数据库并且原生和Kafka Connect整合社区活跃度更高。做CDC链路要注意几个问题。第一Binlog的解析会占用数据库IO高峰期如果没控制好可能影响在线业务一般建议从库提取Binlog。第二CDC消息里一个Update可能带来两条记录变更前和变更后Flink处理时要定义好更新语义否则结果会错。第三数据库表结构变更时CDC消息的Schema会变Flink作业需要兼容老字段和新字段这是线上最容易出事故的环节。4. 实时数仓怎么搭分层、指标和维表设计4.1 实时数仓同样需要分层很多人以为实时数仓就是把离线数仓的表加速一遍这是最大的误解。离线数仓的分层逻辑在实时场景里依然成立但实现方式完全不同。我常用的实时数仓分层模型是四层ODS操作数据层放最原始的实时数据比如Kafka里的订单原始消息这个历史数据一般只保留一两天DWD明细数据层做清洗和标准化把订单流和支付流Join成一条事实明细同时保证字段口径一致DWS汇总数据层做轻度聚合把按分钟、按小时聚合的结果算好存下来比如各店铺每分钟的GMV不对用户开放查询ADS应用数据层面向具体的业务场景比如大屏、报表、个性化推荐数据量小、查询效率高。这个结构跟离线数仓很像但区别在于实时场景每一层都是持续运行的。ODS层的Kafka topic有生命周期DWD层的Join是实时流关联DWS层的聚合结果是增量更新的ADS层则是持续写入Redis或者ClickHouse。生产实践里DWD和DWS之间的口径一致性是最需要花时间对齐的地方因为流式计算是增量的一旦口径变化改起来比离线麻烦得多。4.2 实时指标加工要特别小心幂等性在线指标里最常见的是GMV、UV、PV、订单量、退款率。这些指标在离线里很容易算但在实时里要保证数据重放一遍结果不变这就是幂等性。GMV的加工逻辑通常是支付成功事件加一笔金额退款事件减一笔。这里要防止的问题是同一个支付事件在Kafka重平衡、Flink重启后被重复处理导致GMV算多了。方案有两种一是用订单ID做去重状态Flink里用状态存最近N小时的订单ID重复消息直接丢弃二是让Sink支持幂等更新比如ClickHouse的ReplacingMergeTree用事件ID做去重键后到的重复数据覆盖先到的。UV计算更麻烦因为需要去重。实时UV的经典做法是使用HyperLogLog这样的近似去重算法它能用很少的内存估算出去重后的数量误差通常在0.5%以内业务上是能接受的。如果要求精确去重就得在状态里维护所有访问用户的ID集合数据量一大就扛不住。选哪种方案取决于这个UV是给媒体宣传用的好看指标还是给精准决策用的可靠指标前者近似即可后者需要精确但成本也高。4.3 维表关联订单流离不开商品和用户信息订单流里通常只有商品ID、用户ID、店铺ID这些外键但看板要展示的是商品名称、店铺名称、用户等级。这就需要实时把事实流和维表做关联。最简单的办法是查询外部数据库比如每来一条订单就去MySQL查一次用户和商品信息。这个方案在低流量下能用但大促时会把数据库查死。常见的优化思路有三种一是维表缓存利用Flink的异步IO和本地缓存把热数据缓存在算子内部减少数据库查询次数但缓存有过期时间可能查到旧的维度值二是广播维表把几万条以内的商品/店铺维表做成广播流每个并行的Task都持有一份全量维表关联直接在本地完成三是使用Flink SQL的Temporal Join它根据事件时间自动关联那个时间点上的维度版本这对维表变化频繁的场景很有用。我强烈建议能用广播维表解决的就用广播维表简单粗暴有效。只有维表数据量特别大、广播内存不够时才退回去做异步查询加缓存。4.4 数据质量实时同样需要对账实时的数据质量通常比离线更难保障因为数据是边算边出的没法事后统一修正。我的经验是必须建立三条防线。第一条是延迟监控。Kafka消费Lag、Flink作业处理延迟、结果写入延迟都要设置监控和告警。大屏上数字超过三分钟不动比数字错了更敏感因为业务方会直接看到。第二条是丢失监控。日志采集端要记录发送量和成功量Kafka端要做消息条数统计Flink处理端也要做输入输出条数对比三个数字对不上就说明中间有丢失。第三条是周期对账。离线数仓每天跑完后可以拿离线汇总和实时累计做对比差太多就说明实时链路有逻辑问题。对账的周期不用太频繁一小时一次或者一天一次都行但必须跑。没有对账机制的实时链路就是在裸奔。5. 生产稳定性背压、数据倾斜、重启恢复一次说清5.1 背压是链路健康的晴雨表Flink作业在运行中经常出现一种奇妙现象Kafka里的数据积压越来越多但Flink集群的CPU、内存看起来都不高。其实这是典型的下游处理能力不足导致背压传导某个算子的处理速度跟不上数据到达速度数据在算子内部堆积反压会一级一级往上传最后传到Source端表现为Kafka消费速度暴跌。定位背压的方法很直接在Flink Web UI里看每个算子的BackPressure状态找到压力呈红色的那个算子。常见的背压原因有三类一是单条数据处理逻辑太重比如每条都查外部数据库这种就加缓存或者改为批量处理二是并行度过低分区数大于并行度导致部分Task处理不过来这种就提高并行度三是Sink写入下游太慢比如批量写ClickHouse的批次大小没调好这种就优化Sink的攒批策略。背压不可怕可怕的是链路背压了却没人发现。生产环境一定要对每个作业做Lag监控并且设置分级告警Lag超过阈值就往钉钉群里发警告超过严重阈值就电话通知责任人。我经历过不止一次因为背压导致的大屏数据卡死每次都是第一时间看Lag告警才定位到的。5.2 数据倾斜大促时的头号杀手电商大促时数据天然倾斜爆款商品的订单量可能占全站的三成。Flink做KeyBy聚合时如果按商品ID分key一个热点商品就会把数据全压到一个Task上其他Task闲着这个Task被压垮。解决数据倾斜的常规思路有两种。一种是两阶段聚合先给Key加随机前缀把同一个热点Key拆成多个子Key分别做局部聚合再去除前缀做全局聚合。比如统计店铺GMV时第一步先按店铺ID加随机后缀聚合一部分第二步再把结果汇总。这个方案能有效缓解热点但拆Key会增加一次Shuffle和二次聚合开销需要权衡。另一种思路是用热点检测方案在Flink里实时统计每个Key的流量发现某几个Key异常高时把热点数据单独走一条偏斜处理路径非热点走正常路径。这个方案效果更好但复杂度高适合那种热点极其集中的场景。我建议普通业务先用两阶段聚合实现简单能解决80%的问题。剩下那20%极端热点的场景再考虑动态拆分。不要一上来就搞复杂方案否则维护成本会反过来拖垮团队。5.3 作业升级与状态恢复Flink作业的版升级难点都在状态兼容上。如果作业改了算子逻辑、换了窗口参数状态结构变了从旧checkpoint恢复就可能失败或者状态错乱。规范做法是利用Savepoint做可控的版本升级。Savepoint是手动触发的状态快照升级前先停掉旧作业并触发Savepoint升级代码后再从Savepoint恢复新作业。关键是在修改代码时要理解状态算子Keyed State、Window State的变化凡是删掉或改了状态结构的算子都要评估兼容性。另一个常见的坑是并行度调整。Flink 1.15之后支持从Checkpoint恢复时调整部分算子的并行度但Keyed State的正常恢复还是依赖Key的分布。调整并行度后要注意状态是否能正确重新分配给其他Task否则会出现恢复后数据错乱。生产环境的经验是大版本升级之前一定要在测试环境完整演练一遍从旧作业到新作业的切换流程别等线上出事了才想起验证。5.4 大促容量规划与压测大促的实时链路容量规划做不好就是事故。我一般按三步来做首先是估算峰值流量。从历史大促数据看峰值TPS和平时均值的倍率再结合当年活动的预期增幅给每层链路算出一个目标值。比如平时订单消息5000 TPS大促倍数20倍目标就要做到10万 TPS再乘1.5的冗余系数。然后是压测。大促前一个月就要做全链路压测不只是压Flink作业而是从日志产生、Kafka写入、Flink消费到ClickHouse写入全部压一遍。压测过程中要特别关注Kafka的分区数是否够消费者并行度Flink的Slot资源是否足够下游ClickHouse的写入并发有没有触顶。最后是降级预案。不管怎么压测线上都可能出意外。实时架构一定要保留降级能力最理想的是实时链路异常时大屏可以切换到离线数据的伪实时版本先保证页面有数再慢慢修实时链路。很多团队忽略降级方案结果大促零点实时链路一断全公司的人对着一个静止的大屏发呆那是相当难受的。6. 实战案例从订单库到GMV实时看板的完整链路6.1 场景和指标定义假设现在要做一个支付GMV实时看板展示今日累计GMV、过去5分钟支付订单量、各商品类目的实时排名。数据源是订单系统的MySQL库业务方要求秒级延迟允许大促期间有30秒以内的波动。指标口径需要提前定清楚GMV只算用户实际支付成功的金额不包括退款退款事件要实时扣减已下单未支付的不算。这些口径如果业务方没定清楚开发过程中一定会反复改。6.2 采集与传输层实现订单库的Binlog用Canal监听解析出支付成功和退款两类事件转成JSON消息写入Kafka的topicdwd_order_paid和dwd_order_refund。Kafka分区数设置为24这样下游Flink最多可以开到24个并行度消费。主题里的消息格式用统一的JSON比如支付事件{ order_id: 202412001234, user_id: 998877, item_id: 500123, shop_id: 30001, cat_id: 1201, pay_amount: 199.00, pay_time: 2024-12-01 10:30:00 }这里日期就用通用日期业务跑起来要看自己的数据情况。Canal启动后监控消费Lag确保Binlog的解析速度跟得上写入速度。这块是链路的地基地基不稳后面全白搭。6.3 Flink SQL 主计算流程Flink从Kafka消费消息首先做清洗过滤掉测试订单和金额异常的数据然后分别处理支付和退款事件最后按窗口聚合。用Flink SQL实现先用DDL定义Kafka的Source表CREATE TABLE kafka_order_paid ( order_id BIGINT, user_id BIGINT, item_id BIGINT, shop_id BIGINT, cat_id BIGINT, pay_amount DECIMAL(10,2), pay_time TIMESTAMP(3), WATERMARK FOR pay_time AS pay_time - INTERVAL 30 SECOND ) WITH ( connector kafka, topic dwd_order_paid, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, format json );然后分钟级聚合结果写入Doris供大屏查询INSERT INTO ads_gmv_1min SELECT cat_id, TUMBLE_START(pay_time, INTERVAL 1 MINUTE) AS stat_time, COUNT(DISTINCT order_id) AS order_count, SUM(pay_amount) AS gmv FROM kafka_order_paid GROUP BY cat_id, TUMBLE(pay_time, INTERVAL 1 MINUTE);再维护一个今日累计GMV的指标这个不能直接基于分钟结果累加因为退款要实时扣减。最稳妥的方法是在Flink里维护订单维度的状态记录每个订单当前的净金额再周期性把所有订单金额累加起来。这个逻辑也可以用SQL中的带状态聚合实现核心是更新和删除事件要正确处理。6.4 上线后的效果与踩坑记录这套链路跑起来之后大屏的指标延迟基本稳定在5到15秒。最初踩了几个坑我记下来给大家参考。第一个坑是Canal的Binlog点位重置。Canal重启后会默认从最新位点开始消费如果不是从指定位点启动会漏掉重启时间段内的数据。解决方式是把位点固化到ZooKeeper或者内部存储中每次重启从上次的位置继续。第二个坑是ClickHouse的写入抖动。批量写入时如果某一批数据过大或者集群刚好在做Merge写入延迟会突然拉高导致Flink的Sink算子背压。后来在Sink端做了攒批大小的限制并设置了写入超时和重试策略抖动明显减少。第三个坑是Flink作业更新时的数据双跑。每次改SQL逻辑旧作业和新作业同时跑一段时间结果会出现大屏数字跳变。现在的做法是切换之前先停旧作业、做Savepoint新作业从Savepoint恢复同时旧Sink表不写入等新作业延迟追平后再切换。确保同一时刻只有一套计算在写结果。这套订单实时统计链路从架构上看并不复杂但每一层都有自己需要注意的点。电商数据实时处理架构的本质就是把数据产生、传输、计算、存储、查询这条链路用可靠的方式串起来同时保证延迟、准确性和成本在一个可接受的范围内。没有一套架构能适配所有业务但把分层逻辑、选型依据和常见坑理解透你就能针对自己的场景设计出足够稳的方案。最后分享一个实际操作中很有用的检查习惯每两周做一次全链路的Lag和数据量核对把Kafka的生产消费条数、Flink的输入输出条数、ClickHouse的写入条数放在一张表里对比。这个习惯帮我发现过好几次隐蔽丢数据的问题。实时架构稳定运行靠的不是一次大促的临时保障而是平时持续的体检和冗余建设。
返回列表