
做电商实时数据最难的不是写代码而是把整个架构的“故事”想清楚。我在这行摸爬滚打了十几年从早期的 T1 离线报表到后来的 Lambda 架构再到现在的实时数仓踩过的坑可以写一本书。今天这篇东西我不讲空泛的理论就围绕着“电商数据实时处理架构”这件事把我在实际项目中怎么设计链路、怎么选型、怎么调优、怎么排查问题从头到尾捋一遍。如果你正准备动手搭一套实时数据处理体系或者已经在用 Flink、Kafka 但总觉得哪里不对劲这篇文章应该能给你一些实实在在的参考。1. 电商实时处理架构别急着选型先把故事讲清楚很多人一上来就问“用 Flink 还是 Spark Streaming”“Kafka 要不要上集群”但我建议先停一下。实时处理架构的本质是把“数据从产生到被业务消费”的时间窗口压缩到分钟级甚至秒级。在电商场景里这个时间窗口决定了你能做什么、不能做什么。1.1 实时和离线的定位完全不同离线数仓解决的是“昨天发生了什么”实时数仓解决的是“现在正在发生什么”。这两者的技术路线、数据模型、容错策略差别非常大。我见过不少团队直接把离线那套 HiveSpark 的模型搬到实时里结果就是延迟降下来了但数据不准了或者运维复杂度爆炸。根子在于离线的核心假设是“数据是完整的、可回溯的”而实时的核心假设是“数据是流式的、有延迟的、可能乱序的”。这两个假设不换过来后面所有设计都会拧巴。拿最经典的订单统计来说。离线任务在凌晨跑把前一天所有订单汇总一遍结果准、成本低但只能看“昨天”。实时任务要的是“此刻的 GMV”“此刻的订单量”数据一条条进来每一条都可能晚到、乱序、重复。这就要引入水位线Watermark、窗口、状态后端、精确一次语义这些概念。如果对这些东西没有敬畏心后面必然出问题。1.2 电商场景里典型的实时需求长什么样我把这些年遇到的实时需求归了归类大致是这五类实时大屏与经营看板大促期间的大屏每秒钟都要刷新 GMV、订单量、支付转化率这是最经典的实时场景。这类需求对延迟敏感对准确性要求也很高数字跳错了领导会直接看到。实时风控与反欺诈比如同一设备短时间大量下单、频繁修改收货地址、支付失败后重试异常等都需要在秒级识别并拦截。这里对延迟的要求比大屏更苛刻而且要能回溯会话上下文。实时个性化推荐用户刚浏览了某商品下一秒推荐流就该出现同类商品。这要求实时行为数据能快速进入特征计算和离线训练好的模型一起做推理。实时库存扣减与超卖防控秒杀场景下的库存扣减既要准又要快。这其实是个典型的分布式事务问题但实时数据链路在其中扮演了很重要的角色比如把库存流水实时汇总到风控和调度中心。实时运营触达用户加购但没下单、领了券没用、直播里点了链接没付款这些行为都要在短时间内触发短信、Push、优惠券。实时计算不光要算指标还要能产出“动作”。这些场景看起来杂但背后的架构诉求是一致的低延迟、高吞吐、可容错、数据一致。后面所有技术选型和方案设计都是围绕这四个词展开的。2. 链路设计一条订单数据从产生到可用的完整流转实时数据处理链路从宏观上看就三段数据进得来、算得动、出得去。但每一段里面都有很多讲究。2.1 数据采集层消息队列选型不能拍脑袋绝大多数电商系统的业务数据都在 MySQL、PostgreSQL 这类关系型数据库里也有不少在埋点日志里。要把这些数据实时地搬出来第一步就是选一个可靠的消息队列。Kafka 是这个领域的事实标准但“用 Kafka”和“用好 Kafka”是两回事。我在项目里常用的采集方式是Canal KafkaCanal 伪装成 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 的实时数仓最深的感受是实时架构不是买几个组件搭起来就完事了它是一个需要持续迭代、持续打磨的系统工程。每当你觉得链路已经稳定了大促一来新的瓶颈就会出现——但这也正是这个领域最让人着迷的地方。如果你正在搭或者准备搭一套实时数据处理架构我的建议是先花时间把需求场景理清楚把分层模型规划好再去选型和编码。路走对了后面才能走得远。