十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

电商实时数据管道实战:从Flink流处理到实时数仓

电商实时数据管道实战:从Flink流处理到实时数仓 在做电商数据平台的这几年我最常被问到的一句话是“现在的实时销量到底是多少”以前我只能回答跑批任务得等到明早六点。这种话术在大促场景里几乎等于承认业务失控。后来我们开始认真碰大数据流处理才慢慢摸清楚它在电商场景里到底应该怎么用、用在哪、用到什么程度。这篇文章不打算复述那些引擎官方文档更多是聊聊落地过程中那些真实的需求、真实的选择和真实的坑。如果你也在做电商相关的大数据平台或者在考虑用流处理解决业务实时性问题这篇应该能给你提供一条相对清晰的参照路径。1. 电商场景里流处理到底被什么逼出来的——不是技术视角是业务视角1.1 从一次促销超卖说起实时数据缺口是怎么暴露的事情要从一次常规的秒杀活动说起。当时的架构并不算落后订单库用MySQL分库分表库存扣减走Redis Lua脚本数据分析走Hive离线数仓。活动前压测数据都挺漂亮大家信心满满。结果活动开始没多久运营在后台看到实时销量只有实际订单量的三分之一而库存却已经告急。两边一核对发现同一个SKU的销量和库存根本对不上运营直接慌了临时下线了几个爆款SKU。回过去看问题出在哪订单数据实时写入MySQL但负责统计销量的大屏读的是离线数仓离线任务每隔半小时才刷一次遇上大促瞬时流量半小时的数据差距可能是几十万单。更麻烦的是运营为了补数据临时写了一批脚本直连数据库查询把核心订单库的查询压力直接拉满整个交易链路都开始抖动。这次事故让我彻底意识到一个问题电商业务天然是实时的但技术架构却还停留在批处理思维里。库存扣减是实时的用户支付是实时的优惠券核销是实时的可所有分析和决策链路全是隔了半小时以上的离线逻辑。这个割裂就是流处理要解决的核心问题——不是要炫技术而是业务已经等不起批处理的下一次调度了。1.2 电商链路中哪些环节是真正强实时的聊流处理之前先得把电商链路拆开。不是所有环节都需要毫秒级响应也不存在一个万能的实时方案。我按业务对实时性的敏感程度把电商链路分成了几个梯度。第一梯度是库存和风控。库存扣减要求秒级甚至毫秒级准确超卖就是资损不需要解释。风控同样强实时羊毛党从注册到下单可能只需要十几秒等离线特征跑出来钱已经被刷走了。第二梯度是运营决策和用户体验。实时销量、实时销售额、实时转化率这些是运营随时在盯的指标。大促期间运营要根据实时销量判断是否补货、是否调整投放策略延迟半小时和没有数据没什么区别。再看用户体验侧比如购物车数量的实时同步、物流轨迹的实时更新用户感知极强延迟太久用户就投诉了。第三梯度是离线分析可以覆盖的后验场景。比如月度经营分析、用户画像建模、品类销售趋势这些不需要实时批处理反而更合适因为计算成本低、回溯方便。很多人以为流处理就是把批处理加速其实不对。流处理解决的是第一和第二梯度的问题它和批处理是互补关系。如果一个场景只需要明天知道结果千万别强行上流计算那是在给团队找麻烦。1.3 实时的成本边界不是所有实时都划算实时是有成本的而且是显性成本。流处理作业需要常驻计算资源Kafka集群要持续扛写入流量状态后端要占内存或磁盘监控告警又得多一套体系。这些成本在大促场景下能换来真金白银但在日常场景里就得掂量掂量。我见过一个团队把所有离线报表都改成实时任务理由是“老板要看实时数”。结果集群成本翻了快三倍而老板真正实时看的其实就三张屏。剩下那些明细报表用户点开频率极低实时价值和成本完全不成正比。所以做技术选型前一定要先回答一个问题这个指标如果在五分钟后才更新业务会出什么事如果不会出大事就别走流处理。把流处理留给真正能创造业务价值、避免资损的场景这是电商流处理落地最重要的一条边界原则。2. 流处理引擎选型与架构设计电商团队需要避开的决策误区2.1 Flink、Kafka Streams、Spark Streaming选型背后的取舍引擎选型是整个流处理落地里最容易陷入争论的环节。我在不同阶段用过Spark Streaming、Kafka Streams和Flink感受很不一样。早期我们用Spark Streaming做实时统计它和Spark生态打通上手成本低但微批模型在处理事件时间和精确一次语义时非常别扭。尤其是大促高峰期微批间隔稍微调大指标延迟就明显调小了集群资源压力又上来。后来换成Flink最直观的感受是延迟从秒级变成了毫秒到秒级事件时间、水位线、状态管理这些机制都是原生支持的处理乱序数据不用再自己搞一套复杂的补偿逻辑。Kafka Streams是另一种思路。如果你的实时链路简单不需要复杂的状态计算和窗口聚合直接用Kafka Streams嵌在应用里最轻量。它没有独立的计算集群部署和管理成本都很低。但一旦涉及复杂的多流join、大状态管理和精确一次端到端保证它就不如Flink成熟。我的建议是电商这种业务复杂度下中小团队直接上Flink别犹豫。不是因为Flink最流行而是因为电商流处理绕不开状态、窗口、乱序和端到端一致性这四个难题Flink在这四个维度上的生态是最完整的。Kafka Streams可以作为轻量级数据管道组件存在但要支撑电商全链路实时业务还是Flink更稳。2.2 事件时间、水位线与乱序为什么GMV统计必须关注这三件事电商数据有一个天然属性用户行为发生的时间点和数据真正到达计算引擎的时间点往往不一致。用户点击下单之后前端可能先做了本地缓存或者用户网络不好请求延迟了几十秒才发出去。如果按数据到达时间做统计GMV指标就会出现明显的波动——流量高峰时数据迟到严重统计结果偏低低谷时反而补上了前面的单量指标忽高忽低运营根本没法看。Flink的事件时间机制就是解决这个问题的。它让计算基于业务事件真正发生的时间而不是数据到达的时间。但问题来了事件时间必须要配合水位线一起用水位线本质上是告诉Flink“到这个时间点为止我估计后续的数据已经基本都到了可以把窗口触发计算了。”水位线设多长直接决定了统计的准确性和实时性之间的平衡。设太短大量迟到数据会被丢到窗口之外设太长指标延迟变大。我处理过一条埋点链路移动端的网络环境复杂数据迟到严重我们最后把水位线设到30秒同时开启迟到数据重算机制才让实时GMV和离线GMV在账面上基本对齐。这个数值不是拍脑袋定的是根据线上数据迟到分布统计出来的大家做的时候一定要用自己的真实数据去测别用默认参数糊弄。2.3 精确一次语义概念很美好落地要看清代价精确一次是什么意思简单说就是每条数据无论处理多少次最终结果都和只处理一次一样。电商场景里这个需求非常硬核比如订单金额统计如果数据在计算过程中被重复计算那么实时GMV就是错的财务核对时根本解释不清。Flink checkpoint机制给了精确一次的能力但这里面有个关键容易被忽略checkpoint只能保证Flink内部状态的精确一次如果数据已经写到了外部系统比如MySQL、ES、Kafka外部系统的写入是重复的还是精确的需要靠外部系统的幂等性或事务机制来保证。也就是说端到端精确一次光靠Flink还不够。我见到最常见的问题是把数据从Kafka消费出来经过Flink计算后写到另一个Kafka或者ES写的时候用了普通的producer没有做幂等设计。一旦Flink任务重启并从头恢复状态下游就会收到重复数据。要解决这个问题要么下游写入做去重要么用Flink的两阶段提交机制配合支持事务的sink要么在业务数据里带上全局唯一ID写一张去重表。这里面没有银弹每个方案都要看你的下游存储和业务容忍度。2.4 分层架构和集群部署策略的常见形态等到引擎选型定了紧接着就是架构怎么搭、集群怎么部署。这不是简单地开几个Flink任务就完事而是要有一套相对稳定的分层结构。我当前比较认可的电商实时架构是这样分层的接入层负责对接业务数据源包括订单库的binlog、用户行为埋点、库存变动消息统一写入Kafka。Kafka在这里承担缓冲池的角色削峰填谷下游处理系统的消费速度可以独立调控。计算层是流处理的核心包括实时ETL、指标聚合、规则引擎。这一层用Flink跑按业务域拆分成多个作业作业之间通过Kafka解耦。比如订单域作业把订单信息清洗成标准格式写入实时订单明细Topic下游的GMV统计作业再从这个Topic消费各自算各自的指标。存储层根据查询场景选择不同组件。实时指标用Redis明细查询用ClickHouse或者是Elasticsearch实时大屏直接查预聚合结果避免对底层计算造成压力。集群部署方面电商大促场景下我倾向用Flink on Kubernetes因为弹性伸缩能力更好大促前可以提前扩容大促结束后缩容省成本。如果公司已经有比较成熟的YARN体系Flink on YARN也完全没问题关键是和现有基础设施的整合成本要低。部署时有个点特别容易踩坑并行度不是越大越好。并行度设太高Kafka分区数跟不上会造成部分Task空闲资源浪费并行度太低数据又容易堆积。一般经验是并行度和Kafka分区数保持一致最多是分区数的整数倍同时要给状态访问留出足够的堆外内存。3. 流计算在电商落地中那些我踩过也最常被问的坑3.1 热SKU倾斜问题现象到根因定位数据倾斜是流计算里最常见但也最隐蔽的坑。电商场景尤其明显因为流量天然高度集中在少量爆款SKU上。具体表现是什么Flink作业整体吞吐不高但某个TaskManager的CPU和内存长期跑满其他TaskManager却很空闲。更典型的是Kafka的某个分区写入量明显高于其他分区消费端处理不过来导致整个作业背压。我记得有一次排查实时订单统计延迟发现一个极其热门的商品贡献了全站百分之十几的订单量但Flink按商品ID做keyBy之后这些订单全部压到同一个subtask上单个subtask的处理能力成了瓶颈。怎么解决核心思路是拆分热键。一种方案是给热键的key加上随机后缀把数据打散到多个subtask再在结果合并时二次聚合。比如订单明细按SKU统计时先对每个SKU加10以内的随机后缀拆成10个子key并行统计最终结果再把10个子key的结果合并起来。麻烦的是需要维护热键名单定期更新。另一种方案是利用Flink的rebalance或rescale机制在数据倾斜严重时通过人工改变并行度分配来缓解但这种属于短期治标。3.2 UV去重Bitmap、布隆过滤器还是状态存储电商实时场景中UV独立访客数是最常见也最折腾人的指标之一。PV好算累加就行UV要求去重而去重在大数据量下是非常消耗资源的。最早我们直接在Flink里用MapState存userId逻辑简单但数据量大了就扛不住。一个热门活动页可能要存几百万个userIdRocksDB状态膨胀得厉害作业内存和磁盘占用双双起飞checkpoint时间也越来越长。后来尝试了布隆过滤器内存占用确实降下来了但误差问题让运营非常头疼——明明是新访客被判成老访客UV就虚低。再后来又试过HyperLogLog内存比布隆过滤器更省但误差更大适合做估算而不是精确统计。最终采用的方案是针对精确UV用RoaringBitmap存userId的序号映射对数值型ID压缩效果极好几百万个ID只占几十MB内存。如果量更大可以把用户ID分段存储用多个bitmap组合。这里有一个关键前提——需要给每个用户生成一个连续的数值编号这本身需要一个实时ID映射服务。如果没有这个映射布隆过滤器或者HLL可能就是你当前复杂度下最合适的选择。做技术方案前一定先把“精确”和“近似”这两个选项摆到业务面前让业务来决定而不是默默替业务扛下所有成本。3.3 背压与检查点大促稳定性的两个主要来源大促前我们最担心两件事背压和checkpoint失败。这两件事往往还是连锁反应。背压的本质是上游数据生产速度大于下游数据处理速度导致整个链路被堵住。Flink UI上的背压提示会先从最下游作业开始变红然后逐步向上游蔓延。排查背压第一步是看瓶颈在source、operator还是sink。如果是sink慢比如写入MySQL连接数不够那上游计算再快也没用如果是某个operator慢多半是数据倾斜或者计算逻辑太重。有一次我们的实时GMV作业在大促高峰期背压严重查下来发现瓶颈不在计算而在下游的Redis写入——每次聚合结果都实时写入Redis高峰期写入QPS直接打满写到后面开始出现大量的超时重试反而更堵。最终方案是增加一个批式攒写层窗口结果先攒在内存里每两秒批量写一次Redis压力瞬间降下来了。checkpoint超时也是大促常见问题。Flink定期对状态做快照如果状态太大或者Barrier在链路上传递太慢checkpoint就会超时。状态太大通常意味着状态设计不合理——比如把该放外部存储的大表放进了Flink状态里。如果必须用大状态建议上RocksDB状态后端同时加大checkpoint的并发度和超时时间并根据业务容忍度调整checkpoint间隔。3.4 埋点日志乱序和迟到水位线到底要怎样设置埋点日志的乱序问题比订单数据严重得多。因为埋点数据要经过客户端本地缓存、CDN、网关等多层链路一批数据里既有刚刚产生的最新事件也有十几分钟前的历史事件。如果水位线设得不好实时转化率数字就跟坐过山车似的来回跳动。我处理这类问题的经验是先把数据分布摸清楚再定水位线。具体做法是取最近一周的历史埋点日志统计每个事件从发生到到达Kafka的延迟分布然后取P95或P99作为水位线设定的参考值。比如数据显示95%的数据能在20秒内到达那水位线就设20秒左右再开启迟到数据重算机制处理剩余5%的迟到数据。还有一个小细节实时数仓的表结构设计时事件时间和处理时间一定要分开字段存储不要混着用。事件时间是业务逻辑的基准处理时间是监控排障的依据两个字段都保留后面排查问题会省很多力气。4. 一条可复用的实时指标管道从埋点到实时大屏的完整记录4.1 需求拆解老板要的“实时”到底指什么做实时指标平台最大的坑往往不是技术实现而是需求方自己也没想清楚“实时”是什么含义。我遇到过一个需求运营说要“实时库存”等到我们做完实时链路才发现他们要的其实是每隔五分钟刷新一次的可售库存而不是真正逐笔扣减的实时库存。所以拿到需求之后第一步应该拆解几个问题这个指标多久更新一次能接受数据允许有误差吗历史数据要回溯吗数据要是出了问题影响谁的决策这是做实时平台最重要也最容易跳过的一步但跳过的后果是后期反复改需求。我见过太多团队上来就搭Kafka、写Flink SQL等到大屏上线才发现指标口径和业务认知对不上。口径问题在实时链路里尤其难改因为流计算一旦跑起来改逻辑意味着状态要清空重算不只是改一行代码的事。4.2 管道设计埋点、Kafka分区、Flink作业、ClickHouse出口以我们后来做的一套实时经营指标管道为例完整链路大概是这样的用户行为埋点由前端SDK上报到网关网关把日志写入KafkaKafka按业务事件类型建Topic比如OrderCreateTopic、OrderPayTopic、UserActionTopic。分区的设计要提前想清楚一般按用户ID或者订单ID做key保证同一个用户或同一笔订单的关联事件发到同一个分区这样Flink做基于用户或订单的聚合时单分区内数据是自洽的可以避免跨分区join的麻烦。Flink作业从Kafka消费数据后一部分做基础清洗比如补全订单中的类目、品牌、城市等维度字段一部分做实时聚合比如每分钟算一次总GMV、各品类GMV、Top商品排行。聚合结果有两份出口一份写入Redis供大屏和线上接口查询一份明细数据写入ClickHouse供运营自助分析。这里有个问题值得提醒不要把所有数据都往ClickHouse灌。ClickHouse的单表查询能力很强但写入压力过大时也扛不住而且明细表动辄几十亿行查询性能和成本都会受到影响。一般做法是只把需要明细分析的流量写入ClickHouse其他场景走实时聚合结果。4.3 Flink SQL示例实时GMV、实时库存、实时Top N用Flink SQL写实时指标比用DataStream API开发效率高很多维护成本也低。这里给出一段我们实际用过的简化版的实时GMV统计SQL覆盖普通读者可以复现的主干逻辑CREATE TABLE order_pay ( order_id STRING, user_id BIGINT, sku_id BIGINT, category_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 order-pay, properties.bootstrap.servers kafka-cluster:9092, properties.group.id realtime-gmv-group, format json, scan.startup.mode latest-offset ); CREATE TABLE gmv_minute_sink ( window_start TIMESTAMP(3), category_id BIGINT, gmv_amount DECIMAL(14,2), order_cnt BIGINT ) WITH ( connector jdbc, url jdbc:mysql://dwd-server:3306/realtime_dw, table-name dws_gmv_minute, username realtime_user, password ****** ); INSERT INTO gmv_minute_sink SELECT TUMBLE_START(pay_time, INTERVAL 1 MINUTE) AS window_start, category_id, SUM(pay_amount) AS gmv_amount, COUNT(DISTINCT order_id) AS order_cnt FROM order_pay GROUP BY TUMBLE(pay_time, INTERVAL 1 MINUTE), category_id;这段SQL干了两件事定义Kafka流和MySQL sink然后每1分钟算一次各品类的GMV和订单量。里面的WATERMARK配置就是前文提到的迟到处理策略——允许30秒的乱序。COUNT(DISTINCT order_id)在数据量不大时是可以直接用的但数据量大时建议改成BitMap方案否则状态膨胀速度会超出预期。实时库存的逻辑会稍微复杂一点因为库存变更和订单支付之间有时序关系。我们当时的做法是建一张库存变更明细流包含sku_id、变更类型充值、扣减、回滚、变更数量、发生时间然后按sku_id做累计聚合得到当前实时库存。这个方案要求上游库存操作必须全链路事件化如果还有系统在直接改库存表实时库存永远会对不上。4.4 大屏只是出口真正的产出是能触达业务的动作很多团队做完实时大屏就以为项目结束了其实只做了一半。大屏上的数字本身不产生价值基于数字做出的决策才产生价值。我曾经接手过另一套大屏项目技术链路完全没毛病数据和离线数仓也对得上但运营反馈说这个大屏“没什么用”。后来和运营深聊才知道他们要的不是一个“数字展示板”而是在GMV异常下滑或者转化率异常波动时能提醒他们问题出在哪里。比如“华东区某品类转化率突然下降了5%”最好还能带上关联的流量渠道和商品列表这样运营拿到通知就能直接去调策略。所以后来我们做的实时大屏不只是一组指标卡片还加了实时异常检测和告警模块。Flink作业里用滑动窗口维护每个指标最近30分钟的正常波动范围当实时值超过阈值时自动生成一条异常事件推送到告警群。运营收到告警后点开详情能看到异常指标关联的维度分解就能快速定位问题。这部分是实时平台真正的增值点也最容易被技术团队忽略。如果你在做电商实时平台建议早期就和业务方对齐“看到指标后要做什么动作”这件事否则实时数据很容易沦为摆设。5. 流批一体和实时数仓电商流处理的下一步怎么走5.1 流批一体的现状不是替代而是统一这两年流批一体的口号很响很多团队一上来就说要搞流批一体数仓。但落到实际我觉得流批一体更多是把计算引擎、元数据和口径统一起来而不是用流处理替代批处理。我们目前的实践是数据源统一接入Kafka一份源数据既供Flink实时计算又同步一份到Hive做离线测试和回溯。两套计算共用同一套口径定义比如GMV的计算逻辑、用户活跃的定义都放到公共的指标层统一管理。这样做的原因是随着业务逐步精细化实时和离线口径经常会有偏差财务核对时最怕的就是两边对不上。Flink本身也支持批模式运行同一个作业但我们没有完全依赖这一点因为实际生产里实时任务和离线任务的资源消耗、调度策略、容错手段差异很大。统一是统一口径和元数据不是强行用一套引擎一套架构打天下。5.2 流处理还能延伸到哪里推荐、A/B实验与动态定价除了大家熟悉的GMV和库存实时监控流处理在电商里还有几个特别有价值的应用方向我感触比较深的是实时推荐和动态定价。推荐系统通常依赖用户的历史行为做离线训练但用户当下的兴趣是实时变化的。一个用户刚搜索了某个品类却没有点击推荐位上的商品这个信号如果等半小时后才进入推荐系统用户可能已经离开了。我们做的实时行为特征管道把用户的浏览、收藏、加购行为实时写入特征服务推荐系统在出结果时可以读取这些实时特征对候选商品进行实时重排。这个改造对点击率的提升非常明显。动态定价也是流处理的好场景。大促期间库存和流速变化极快如果人工作决策根本跟不上。我们曾用流处理实时统计每个SKU的销量、库存、加购转化率等特征喂给定价策略引擎当某个商品库存低于安全水位且流速持续上涨时系统会自动调整促销折扣力度防止爆款快速售罄后无货可卖。这些场景的共同点是它们都严格要求低延迟和连续计算批处理无法完成。流处理在这些领域发挥的价值会逐渐超过最传统的实时报表。5.3 团队引入流处理前想清楚三件事如果你所在团队正准备从0到1搭建流处理能力或者说刚刚踩进这个坑我个人有三条建议。第一先选一个再小不过的真实业务场景跑通全链路。不要一上来就追求大而全的实时数仓找一个一天只有几万条数据的订单状态同步场景把Kafka、Flink、ClickHouse的链路完整跑通让团队所有成员都亲手摸过一遍。第二流处理的运维成本必须在一开始就预算进去。Kafka集群监控、Flink作业管理、checkpoint监控、数据积压告警这些不是可以后补的东西。没有监控就上生产环境等于裸奔。我们早期因为监控缺失一次数据倾斜问题折腾了两天才定位到那两天业务方每天都在催数。第三业务口径的拉通比技术实现重要得多。流处理上线前一定要拉着运营、财务、数据团队坐下来把每一个指标的定义、统计口径、更新频率都确认清楚形成文档。别以为这些能事后补充流计算一旦上线改口径就意味着等状态清理完再从某个时间点重算代价远比批处理大。电商行业的流处理说到底不是为了追逐技术热点而是为了把那些过去藏在批处理调度间隙里的业务机会找回来。每一条实时数据的背后都可能对应着一笔可以挽回的订单、一位不该流失的用户、一次不该发生的超卖。能把这件事做到位比用什么引擎、跑多少并行度都重要。
返回列表