做电商实时数据,最难的不是写代码,而是把整个架构的“故事”想清楚。我在这行摸爬滚打了十几年,从早期的 T+1 离线报表,到后来的 Lambda 架构,再到现在的实时数仓,踩过的坑可以写一本书。今天这篇东西,我不讲空泛的理论,就围绕着“电商数据实时处理架构”这件事,把我在实际项目中怎么设计链路、怎么选型、怎么调优、怎么排查问题,从头到尾捋一遍。如果你正准备动手搭一套实时数据处理体系,或者已经在用 Flink、Kafka 但总觉得哪里不对劲,这篇文章应该能给你一些实实在在的参考。
1. 电商实时处理架构:别急着选型,先把故事讲清楚
很多人一上来就问“用 Flink 还是 Spark Streaming”“Kafka 要不要上集群”,但我建议先停一下。实时处理架构的本质,是把“数据从产生到被业务消费”的时间窗口压缩到分钟级甚至秒级。在电商场景里,这个时间窗口决定了你能做什么、不能做什么。
1.1 实时和离线的定位完全不同
离线数仓解决的是“昨天发生了什么”,实时数仓解决的是“现在正在发生什么”。这两者的技术路线、数据模型、容错策略差别非常大。
我见过不少团队,直接把离线那套 Hive+Spark 的模型搬到实时里,结果就是:延迟降下来了,但数据不准了,或者运维复杂度爆炸。根子在于,离线的核心假设是“数据是完整的、可回溯的”,而实时的核心假设是“数据是流式的、有延迟的、可能乱序的”。这两个假设不换过来,后面所有设计都会拧巴。
拿最经典的订单统计来说。离线任务在凌晨跑,把前一天所有订单汇总一遍,结果准、成本低,但只能看“昨天”。实时任务要的是“此刻的 GMV”“此刻的订单量”,数据一条条进来,每一条都可能晚到、乱序、重复。这就要引入水位线(Watermark)、窗口、状态后端、精确一次语义这些概念。如果对这些东西没有敬畏心,后面必然出问题。
1.2 电商场景里典型的实时需求长什么样
我把这些年遇到的实时需求归了归类,大致是这五类:
- 实时大屏与经营看板:大促期间的大屏,每秒钟都要刷新 GMV、订单量、支付转化率,这是最经典的实时场景。这类需求对延迟敏感,对准确性要求也很高,数字跳错了领导会直接看到。
- 实时风控与反欺诈:比如同一设备短时间大量下单、频繁修改收货地址、支付失败后重试异常等,都需要在秒级识别并拦截。这里对延迟的要求比大屏更苛刻,而且要能回溯会话上下文。
- 实时个性化推荐:用户刚浏览了某商品,下一秒推荐流就该出现同类商品。这要求实时行为数据能快速进入特征计算,和离线训练好的模型一起做推理。
- 实时库存扣减与超卖防控:秒杀场景下的库存扣减,既要准又要快。这其实是个典型的分布式事务问题,但实时数据链路在其中扮演了很重要的角色,比如把库存流水实时汇总到风控和调度中心。
- 实时运营触达:用户加购但没下单、领了券没用、直播里点了链接没付款,这些行为都要在短时间内触发短信、Push、优惠券。实时计算不光要算指标,还要能产出“动作”。
这些场景看起来杂,但背后的架构诉求是一致的:低延迟、高吞吐、可容错、数据一致。后面所有技术选型和方案设计,都是围绕这四个词展开的。
2. 链路设计:一条订单数据从产生到可用的完整流转
实时数据处理链路,从宏观上看就三段:数据进得来、算得动、出得去。但每一段里面都有很多讲究。
2.1 数据采集层:消息队列选型不能拍脑袋
绝大多数电商系统的业务数据都在 MySQL、PostgreSQL 这类关系型数据库里,也有不少在埋点日志里。要把这些数据实时地搬出来,第一步就是选一个可靠的消息队列。
Kafka 是这个领域的事实标准,但“用 Kafka”和“用好 Kafka”是两回事。我在项目里常用的采集方式是Canal + Kafka:Canal 伪装成 MySQL 的从库,订阅 binlog,把增删改都变成消息写到 Kafka。这套方案成熟、坑少,但有几个点要特别注意。
第一,binlog 的格式。MySQL 的 binlog 有三种格式:STATEMENT、ROW、MIXED。Canal 要求必须用 ROW 格式,因为只有 ROW 格式才记录了每行数据的完整变化。如果你线上库还在用 STATEMENT,改的时候一定要评估对现有系统的影响,不少团队在这里翻车。
第二,Topic 的分区规划。Kafka 的吞吐跟分区数直接相关,但分区不是越多越好。我的经验是:分区数 = 消费者线程数 × 单分区预期吞吐。一般单分区能扛 5~10MB/s 的写入,先按未来半年数据量预估,宁可少分,后面再加分区要重新分配数据,很麻烦。
第三,消息的 key 设计。如果按订单 ID 做 key,同一订单的所有变更都会进同一个分区,消费者本地就能保证顺序;如果按用户 ID 做 key,那一个用户的所有行为都串行处理,适合做用户画像。这个选择会影响下游计算逻辑,必须提前定。
第四,日志采集。除了数据库,还有大量前端埋点和服务器日志。这路数据我一般用 Filebeat 或者 Fluentd 采到 Kafka。日志数据量大但单个价值低,可以单独用一套 Topic,设置更短的保留时间,避免占用主链路的磁盘。
2.2 计算层:为什么最终选了 Flink
实时计算引擎这块,市面上主要就是 Flink 和 Spark Streaming。很多老团队从 Spark 起家,到了实时这块自然想复用 Spark 的技术栈。但我自己的经验是,如果做真正的实时处理,Flink 是更顺手的选择。
核心原因有三条。
- 真正的流式处理:Spark Streaming 本质是微批,把数据攒一小段再处理,虽然有 Structured Streaming 做了改进,但仍然是基于批的模型。Flink 是原生的流式引擎,数据一到就处理,延迟能到毫秒级,窗口和事件时间处理也更自然。
- 状态管理能力:实时计算很多时候要“记住”之前的数据,比如双流 Join、去重、窗口聚合。Flink 有内置的状态后端,支持 RocksDB、内存、文件系统多种存储,还能和 Checkpoint 机制结合,做故障恢复。这块 Spark 虽然也做,但不如 Flink 成熟。
- 生态和社区:Flink 在实时领域的社区活跃度、文档、算子丰富度,这几年都明显压过 Spark。你遇到问题搜一圈,基本都有现成答案。
当然,选 Flink 也不等于万事大吉。Flink 的调优复杂度不低,尤其是状态、窗口、Checkpoint 这些核心机制,理解不到位很容易踩坑。这个我在第三部分细讲。
2.3 存储层:实时计算完,数据放哪去
计算引擎算出来的结果,必须写到一个能扛住高并发查询的存储里。这里有个常见的误区:把结果直接写回 MySQL。
MySQL 扛不住大促期间的实时大屏。每秒几千次的写入和查询,MySQL 很快就 CPU 打满,慢查询一堆。我的建议是分场景选存储。
- 实时大屏、即席查询:用ClickHouse或Doris。这两个都是列式存储,聚合查询非常快。ClickHouse 的 MergeTree 引擎配合预聚合,秒级返回上亿数据的聚合结果;Doris 在实时更新和标准 SQL 兼容性上更友好。
- 精确到用户的实时数据:比如用户实时画像、实时推荐特征,用Redis。这类数据要求点查性能极高,而且经常要设置过期时间。
- 实时报表的历史归档:实时算出来的结果,定期同步到离线数仓(比如 Hive 或 Iceberg),用于长期分析和回溯。
存储选型不是越贵越好,核心是匹配查询模式。你要让人查“近 5 分钟的 GMV”,ClickHouse 一个聚合就出来了;你要让人查“某个用户的 30 天行为时间线”,那得靠 Redis 或者 OLTP 库。把查询模式捋清楚,存储自然就选出来了。
3. 核心落地:订单实时统计链路的完整实现
说完了架构,我们用一条最核心的链路——“订单实时统计”——来走一遍完整实现。这条链路几乎每家电商都得做,麻雀虽小五脏俱全。
3.1 实时数仓分层:ODS、DWD、DWS、ADS 各自干啥
不少人觉得,实时不就是写个 Streaming 任务,从 Kafka 消费然后算个总数嘛。这么想的团队,前三个月会很快,后面业务一复杂就崩。我的做法是严格按数仓分层的思路来做实时链路,每层有每层的职责。
- ODS 层(原始数据层):Kafka 里的原始消息,直接映射业务表结构,字段基本不动。这层的作用是保留原始数据,方便回溯和重算。比如订单表 binlog 消息、用户行为日志,都在这一层。
- DWD 层(明细数据层):对流式数据进行清洗、补全、规范化,形成事实明细。比如订单表要关联商品维度表,补上商品类目、店铺名称;行为日志要规范化用户 ID、设备 ID 等字段。
- DWS 层(汇总数据层):按业务主题做预聚合,比如按店铺、按类目、按时段汇总订单金额和数量。这层是实时大屏的主要数据源,也是查询性能的关键。这一层通常会写入 ClickHouse。
- ADS 层(应用数据层):面向具体业务应用,比如大屏展示、告警服务、推荐系统的输入。ADS 可以非常薄,就是从 DWS 查数据再加点业务规则。
分层的最大好处是:当业务方提出一个新指标时,你不用从原始数据重新跑一遍链路,很多公共的明细和汇总已经在 DWS 里了,你只需要新增一个 Flink 任务做二次聚合就行。我见过不分层的团队,每次需求变更都要改最底层的任务,改一次全链路重算一次,苦不堪言。
3.2 窗口、水印、Checkpoint:三个必须把握好的核心机制
Flink 最关键也最容易翻车的三个机制,我一个个说。
窗口(Window)。实时统计“最近 5 分钟订单量”就需要用窗口。Flink 有滚动窗口、滑动窗口、会话窗口三种。电商场景里,滑动窗口最常用,比如“近 5 分钟 GMV”就是每 10 秒滑动一次的 5 分钟窗口。这里要设计好窗口大小和滑动步长,窗口太大数据不及时,窗口太小计算量暴涨。我一般建议窗口 1~5 分钟,滑动 10~30 秒,这样既能保证实时性,又不至于把 CPU 打爆。
水印(Watermark)。这是 Flink 事件时间处理的核心机制,也是新手最容易搞混的概念。简单理解,水印就是“我目前收到的数据里,最大事件时间减去我们允许的迟到时间”。它解决的是数据乱序问题:可能 10:00:05 的数据比 10:00:04 的数据先到,你怎么判定 10:00:04 的数据都到齐了?靠水印。
设置水印时要留一定余量,不能太激进也不能太保守。太激进(比如只留 1 秒)会导致大量数据被判定为迟到,结果不准确;太保守(比如留 1 分钟)会导致结果出来慢不少。我在实战里一般留 5~10 秒,还要配合允许迟到(allowedLateness),让迟到的数据触发一次修正更新。
Checkpoint。这就是 Flink 的容错机制,定期把任务状态存到外部存储(HDFS 或 S3)。如果任务挂了,就从最近一次 Checkpoint 恢复。这里最容易踩的坑是Checkpoint 时间过长,超过设定的超时时间导致任务频繁重启。主要原因是状态太大或者网络抖动。我建议把 Checkpoint 间隔设成 1~3 分钟,同时开启增量 Checkpoint,尤其是用 RocksDB 状态后端时,增量能省大量 IO。
3.3 双流 Join 和状态过期,一不小心就是坑
实时计算里最复杂的一块,是双流 Join。比如订单流和支付流要关联到一起,得到“已支付订单”。两个流数据到达时间不一致,订单先到、支付后到,你得在内存里等支付数据。
这就要用 Flink 的状态管理。一个订单进来后,先写到状态里,等待支付流的关联;支付流到了,查状态里有对应的订单就关联上,没有就也存下来等订单。这个状态下大了,内存扛不住,就要用 RocksDB。
但状态不是无限留的,必须设置状态过期时间(TTL)。我的经验是,业务数据流的关联等待时间一般设定 5~15 分钟。超过这个时间还没关联上,要么这条数据就是孤数据,要么就是业务异常数据,继续占着内存只会拖垮整个任务。
这里有个非常经典的问题:Join 导致的重复数据。比如支付消息在某些场景下会产生多条,如果不做去重,关联出来的结果就会翻倍。解决的办法是在 DWD 层就按业务主键做去重,比如按订单 ID + 业务类型做 key,在 Flink 状态里维护最近处理过的记录,重复的直接丢弃。
4. 性能调优和稳定性保障:实战里踩过的那些雷
架构搭起来容易,但要让它在千万级订单、几十万 QPS 的压力下稳定跑,是真功夫。这一部分我分享几个我亲身踩过的雷和对应的解决办法。
4.1 数据倾斜:这个最常见,乱用 Key 等于自杀
Flink 做聚合时,经常要按 key 分组。如果某个 key 的数据量远大于其他 key,那这个 key 所在的并行子任务就会成为瓶颈,整个任务的吞吐被它拖死。
我遇到过一次很典型的情况:按店铺 ID 统计实时销售额,结果某个头部大主播的店铺数据量是其他店铺的几百倍。所有数据都挤在一个子任务里,CPU 打满,其他子任务在空转。整个任务的延迟从秒级变成了分钟级。
解决思路是两阶段聚合。第一阶段,给 key 加一个随机后缀,把数据分散到不同子任务里做局部聚合;第二阶段,去掉后缀,把局部聚合的结果再汇总。这样大 key 带来的倾斜就被打散了。
但要注意,两阶段聚合只能解决“预聚合型”的需求(比如 count、sum),解决不了需要精确去重或者关联的场景。如果非要用大 key 做关联,那就得考虑对状态做分区拆分,或者用支持热点检测的框架来动态调整并行度,这个复杂度就比较高了。
4.2 消费堆积:不是加机器就一定奏效
Kafka 消费者堆积(Lag)是实时系统最经典的问题。堆积的直接原因是消费速度跟不上生产速度。
很多人的第一反应是“加消费者实例”。但加了实例,如果分区数没变,消费者数量超过分区数,多余的实例会空转,因为 Kafka 同一分区在同一时刻只能被一个消费者消费。所以正确的做法是:先保证 Topic 分区数 ≥ 消费者实例数,再考虑扩容。
如果分区数已经合理,还是堆积,那就得排查下游 Flink 任务的瓶颈了。常见的原因有这么几个:
- 窗口计算太重:特别是滑动窗口,每个事件要更新很多个窗口,计算量成倍增长。优化方式是把长窗口拆成短窗口 + 增量聚合。
- 状态读写太慢:RocksDB 状态后端如果配置不当(比如内存开太小、block 缓存不够),读写性能会很差。可以适当调大 state.backend.rocksdb.memory.managed 和 block cache 大小。
- 数据倾斜:原因就是前面说的问题,倾斜的那个子任务拖慢整体。
有个排查思路很管用:打开 Flink 的 Web UI,看每个子任务的背压指标(BackPressure)。如果某个子任务持续 High,说明它所在的反序列化、计算或状态读写有问题;如果是整条链路都 High,那大概率是 Source 端消费能力跟不上。
4.3 重复消费与数据一致性:End-to-End 的 Exactly-Once
分布式系统里,消息丢失和重复是常态。Kafka 提供了 At-Least-Once 和 Exactly-Once 两种语义,但 Exactly-Once 只保证在一个 Kafka 到 Kafka 的链路里生效。如果下游是 MySQL、ClickHouse,那 Flink 的 Checkpoint 机制只能保证算子状态的一致性,落库那一步还是可能重复。
我在项目里的做法是,不冒险依赖纯粹的 Exactly-Once 落库,而是改造成“幂等写入”。比如写 ClickHouse 时,用 ReplacingMergeTree 引擎,以业务主键为去重键,重复写入也能收敛为一条;写 Redis 时,用 SETNX 或者给 key 加业务唯一 ID,重复设置不会引入脏数据。
这套“幂等 + 外置去重”的方案,在工程上比强行追求 Exactly-Once 更稳妥。要知道,Exactly-Once 的代价是大量的状态存储和协调开销,当吞吐量上来时,性价比会明显下降。
5. 可视化与服务:实时数据如何真正帮到业务
链路通了,数据算出来了,最后一步是让人能看到、用到。这块做得不好,整个实时架构的价值会大打折扣。
5.1 实时大屏查询的缓存设计
很多团队把实时大屏直接接到 Flink 结果表上,前端每秒钟轮询一次 ClickHouse。这种方式在大促期间很容易把数据库打垮。我建议在大屏和数据存储之间,加一层查询缓存。
比较成熟的方案是,Flink 把聚合结果写到 Redis,大屏的服务端从 Redis 读,Redis 里没有的再回源 ClickHouse。Redis 的 QPS 能力比 ClickHouse 高一个数量级,而且对热点 key 的访问非常友好。
这里要设计好缓存更新的策略。我的做法是,Flink 窗口计算的每一条结果更新,都带上一个自增版本号或者时间戳,写入 Redis 时用 HSET,大屏端每次读出来直接覆盖展示。这样既保证数据新鲜,又避免频繁写库造成压力。
5.2 实时数据质量监控
实时数据出错了,比离线数据出错更可怕,因为它是“当下正在发生的错误”。我在实战里专门做了一套实时数据质量监控,用来兜底。
核心思路是对账。拿订单实时统计来说,我同时维护两条链路:一条是 Flink 实时聚合,一条是离线数仓的定时汇总。每隔 5 分钟,实时链路的结果会和离线链路最近一次的全量结果做对比。如果偏差超过阈值(比如 1%),就触发告警。这样能提前发现水印设置太激进、窗口丢失数据、或者状态过期导致的数据缺失。
另外一个很实用的实践是埋点追踪。在 Flink 的算子入口和出口都打点,统计每秒钟进多少条、出多少条、丢弃多少条。当入口数量和出口数量长期不一致时,基本可以断定链路某个环节出了问题,能省去大量排查时间。
6. 最后的一点体会
电商实时数据处理架构,说到底是门实践科学。技术选型很重要,但更重要的是对整个数据链路的理解、对业务场景的判断、以及对各式各样故障的应对能力。我从最早用 CronJob 定时脚本轮询数据库,到现在完整落地 Flink + Kafka + ClickHouse 的实时数仓,最深的感受是:实时架构不是买几个组件搭起来就完事了,它是一个需要持续迭代、持续打磨的系统工程。每当你觉得链路已经稳定了,大促一来,新的瓶颈就会出现——但这也正是这个领域最让人着迷的地方。
如果你正在搭或者准备搭一套实时数据处理架构,我的建议是:先花时间把需求场景理清楚,把分层模型规划好,再去选型和编码。路走对了,后面才能走得远。