
米兰智能零售场景里跑库存与供应链实时优化平台这个项目我前前后后做了大半年踩了不少坑也总结了不少经验。当时接手这个任务的时候客户给的诉求其实特别直白门店SKU数动不动就上万促销活动频繁线上线下的订单汇到一起之后人工Excel排补货根本排不过来。更让人头疼的是总部基于前一天数据的库存报表到了第二天下午看的时候已经完全是另一回事了——热门单品售罄没人知道仓库里滞销品却占了大量现金流。说白了零售行业的利润本来就被库存周转率卡得死死的一旦数据跟不上业务节奏优化就无从谈起。这篇文章就是围绕这个项目来写的。我会从需求边界、架构选型、高并发链路设计、算法模型落地以及线上压测与排查几个维度展开核心解决的是如何在高吞吐、低延迟的数据链路中让库存和供应链决策从T1变成T0这个问题。适合正在做零售数据平台、供应链优化、实时计算相关的工程师和架构师参考也适合那些业务部门提了需求但不知道怎么落地的人作为对照。1. 需求还没动手就得先拆清楚米兰门店的库存问题到底卡在哪里1.1 传统批处理模式的瓶颈与业务侧的真实痛点项目刚启动的头两周我们并没有急着选技术栈而是把所有干系人拉在一起把库存不准确这个问题从业务侧到技术侧逐层拆了一遍。当时米兰几家旗舰店的情况是门店POS系统每15分钟做一次增量同步到中央数据库ERP系统每天晚上11点跑一次全量库存汇总WMS仓库系统则只能做到每小时上报一次出入库记录。三个系统各有各的时间差到了总部BI报表那边库存数据是对不上的而且谁也说不清哪个才是真实库存。深入现场走访之后发现门店理货员每天要拿着手持终端扫货架把实际商品数量录入系统这个动作本质上是在给系统打补丁。高峰期比如周末、打折季、节假日货架上的商品流动极快手动盘点永远跟不上变化速度。更大的问题在供应链侧采购部看到系统里某SKU库存剩下5件于是下单补货但实际上这个SKU在门店后仓还有两箱没拆封——系统根本没把后仓和在架库存分开建模。所以我们前期做的第一件正事就是和业务一起重新定义了实时库存的语义。它不是一个简单的数字而是由在架可售库存、后仓库存、在途库存、锁定库存用户下单未支付、门店预留、团购预留四个维度组成的复合状态。只有当这四类数据在同一个数据口径下闭环流转实时优化才有意义。1.2 实时优化的核心诉求拆解哪些场景真正需要秒级响应把需求从我们要做实时库存细化成具体业务场景是我们避免过度设计的关键一步。经过三轮工作坊我们把需求收敛成了五个核心场景。第一是门店大屏和运营后台的库存可视化管理这一层要求数据延迟控制在10秒以内让店长能实时看到在架商品的库存水位支撑即时调拨决策。第二是促销活动的库存实时扣减比如某品牌在米兰市中心店做限量发售线上小程序抢购和门店POS同时出库必须保证不超卖。第三个场景是缺货预测与自动补货通过分析过去15分钟、1小时、24小时、7天的销售速率变化动态调整补货建议和安全库存阈值。第四是跨店调拨建议当A店某SKU库存见底而B店滞销库存偏高系统要在15分钟内给出调拨建议。第五是供应链KPI的实时监控比如订单满足率、库存周转天数、滞销SKU占比这些指标要从T1变成分钟级刷新。这五个场景对数据延迟的要求差异很大有的10秒内就能接受有的真要达到秒级甚至毫秒级。我们最终确定了分级处理的思路——不是所有数据都走同一条高并发链路而是在接入层就按重要性和时效性打标走不同优先级的处理管道。核心交易类数据POS流水、线上订单用实时管道高优先级主数据变更商品信息、门店档案用准实时管道中优先级像供应商交期这类低频变更数据直接走批量管道T1更新也没问题。2. 架构选型与链路设计为什么这单业务我坚持用Lambda而不是全实时流2.1 架构决策的完整思考过程从Kappa到Lambda的反复拉锯在技术选型上我们内部吵过好几轮。团队里有人倾向于全实时流处理Kappa架构觉得既然都做实时了干脆所有计算都在流里完成批处理层干脆不要了。这个想法理论上很漂亮但落到实际场景里就有问题。零售供应链场景有一个特殊性计算逻辑的复杂度和数据回溯需求极高。比如库存周转率这个指标需要关联采购成本、销售折扣、期初期末库存等多个维度的数据又比如畅销品预测模型需要对过去90天甚至更长的历史数据进行重算验证模型参数是否还适用。如果所有历史数据都塞在Kafka里要用的时候再重放成本和复杂度会非常大。而且一旦业务要调整指标口径这种需求在零售行业几乎每个月都有所有历史数据都要重算一遍——全流式架构在这种场景下会让人崩溃。所以最终我们选择的是改良版Lambda架构批流共用同一套数据源Kafka批处理层用Spark定期重算全量指标实时层用Flink只处理增量窗口内的计算最终在ClickHouse里做实时结果和历史结果的合并对外暴露统一的数据服务API。这套架构的灵活性在于既保证了实时场景的秒级响应又不牺牲复杂分析场景的计算准确性。数据链路整体上可以分为五个环节采集层门店POS、线上订单、WMS、ERP等系统的数据接入→传输层Kafka消息队列→计算层Flink实时计算 Spark批量计算→存储层ClickHouse分析型数据库 Redis缓存 MySQL元数据管理→服务层统一数据API支撑门店大屏、补货系统、调拨系统、管理后台。链路看起来不复杂但真正实现起来每一个环节都有各自的坑。2.2 各层技术选型与替代方案对比以及我踩过的选型坑选型对比往往是项目前期最费时间的一件事我把几个核心组件的对比过程拿出来说说。消息队列我们对比过Kafka、Pulsar和RabbitMQ。RabbitMQ在吞吐量上明显不够看直接淘汰。Pulsar的架构确实先进多层存储、计算存储分离这些特性很吸引人但当时团队对Pulsar的运维经验几乎为零而且米兰那边的公有云资源上部署Pulsar的案例也不多。综合考虑运维成本和社区活跃度最终选了Kafka版本用的3.5。这个决定在后来的压测中证明是合理的Kafka在高吞吐场景下的稳定性经历过大量生产环境的考验相关的监控告警、扩缩容方案都相对成熟。实时计算引擎选了Flink1.17版本。虽然在流处理领域也考虑过Spark Structured Streaming但Flink在低延迟、精确一次Exactly-Once语义和窗口管理上的表现确实更适合我们的场景。尤其是在做会话窗口分析比如用户从浏览到下单的整个链路时Flink的灵活性明显更强。不过Flink 1.17版本的CheckPoint机制在状态比较大时容易出现反压这个我们后面在调优环节处理了很久后面详细说。存储层是争论最多的部分。实时分析选了ClickHouse主要看中它的列式存储和向量化执行引擎在聚合查询上的性能。MySQL继续作为业务元数据存储商品档案、门店信息、供应商主数据Redis用来做热数据的缓存——比如Top100热销SKU的实时销量、库存水位这种需要毫秒级响应的数据每次查询都去ClickHouse跑聚合显然不现实。这里我想单独提一下选型踩过的一个坑最初试图用Elasticsearch同时承担实时查询和聚合分析的功能结果发现库存流水一旦上了千万级别ES的聚合查询性能和内存消耗就很难接受。后来才把聚合分析全部迁移到ClickHouseES只保留搜索场景比如商品名称模糊查询、门店地址检索。所以选型一定要想清楚每个组件职责的边界不要指望一个组件搞定所有事情。2.3 数据模型的统一设计从源头避免每个部门一套口径的乱象零售行业数据项目最常见的失败原因不是技术不行而是口径不统一。同样一个库存财务部关心的是库存金额成本价计算运营部关心的是可售件数采购部关心的则是包括在途在内的总供给量。如果每个部门的数据都是各自取数、各自加工那系统做出来之后大家拿着不同的数字开会争论的焦点永远在谁的数是对的而不是怎么解决问题。所以我们花了将近三周的时间梳理了所有数据域的统一口径。还是以库存举例我们定义了物理库存仓库或门店实际有的数量、可售库存物理库存减去锁定库存、可用库存可售库存加上在途库存减去安全库存阈值这三个核心指标并且在全公司范围内确认了计算逻辑。这个动作很费时间但却是整个项目能否成功的先决条件。数据仓库模型设计上我们采用了Data Vault 2.0的思路把数据分成Hub核心实体如商品、门店、供应商、Link实体之间的关系、Satellite实体的属性与状态变化三层这样既保证了扩展性也方便追溯历史变化。3. 高并发数据处理这条链条每一环都是在跟延迟和资源较劲3.1 采集层编排门店POS、线上订单、WMS多源数据接入的乱局整理数据接入第一个要解决的问题就是来源各异的系统该怎么接。米兰门店的POS系统是Oracle数据库线上订单走的是自研微服务APIWMS是SAPERP是另一个老系统而且这些系统分布在不同的网络环境里有的能直连有的只能通过文件交换或者FTP取数。为了让接入层清爽我们开发了一个统一的采集组件基于Canal Debezium 自研HTTP Sink核心思路是对支持CDC变更数据捕获的系统用Canal监听Binlog不支持的系统就用自研采集器定时拉取增量数据最终统一写入Kafka。这里有个经验值得分享对于Oracle数据库的CDC用Debezium的Oracle Connector时一定要特别注意LogMiner的权限配置稍微配置不对就会导致归档日志读取权限报错而且这个问题在测试环境不容易发现一上生产就被日志刷屏。另外门店网络不稳定、断网重连这些情况要考虑充分。我们为采集端设计了本地文件缓冲Tiered Storage也就是门店边缘节点先写本地文件再异步上传到云端Kafka。这个设计在双十一大促或者周末大促时救了命——门店网络偶尔会被POS流量占满导致数据上传超时本地缓冲能保证数据不丢等网络恢复再继续上传。3.2 Kafka吞吐量与分区策略不能让数据到了门口还排长队Kafka的调优是整个链路中最核心的环节之一。刚开始我们把Kafka部署在默认配置上压测的时候发现吞吐量一直上不去生产者侧的发送延迟抖得很厉害。后来逐层排查定位到几个关键配置项。第一是批次大小和延迟时间的平衡。batch.size默认是16KBlinger.ms默认是0。在高并发场景下如果单条消息很小比如POS流水可能只有几百字节每条消息都立即发送会造成大量小请求浪费网络IO和Broker的CPU。我们调成了batch.size64KBlinger.ms20这样让生产者攒一批再发吞吐量提升了将近3倍代价是单条消息的延迟增加了20毫秒——对于库存这个场景来说完全能接受。第二是分区数的设计。最开始的topic只设置了8个分区在并发量上来之后明显出现分区间数据倾斜。后来我们把核心topic的分区调整到与下游消费者并发度匹配的数量32个分区下游Flink算子并行度设为32并且按照SKU_ID作为分区键进行Hash路由保证同一个SKU的库存变更消息始终进入同一个分区从而保证处理顺序。第三是Broker端的acks设置。为了保证不丢数据必须设置acksall所有ISR副本确认但这会带来一定延迟。我们的实际配置是acksallmin.insync.replicas2 3副本这样既保证了数据安全又不会让性能掉太多。压测下来这套配置在3台Broker的集群上单topic可以达到每秒约8万条消息的稳定吞吐。3.3 Flink实时计算链路从Source到Sink的完整实现与优化实时计算的核心逻辑分为三大块订单与库存扣减流、销售速率统计流、补货触发流。这里给出一段Flink SQL的示例展示我们如何实现分钟级的销售速率统计。-- 定义Kafka Source CREATE TABLE pos_sales ( event_time TIMESTAMP(3) METADATA FROM timestamp, shop_id STRING, sku_id STRING, quantity INT, sale_amount DECIMAL(10, 2), WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic pos_sales, properties.bootstrap.servers kafka-1:9092,kafka-2:9092,kafka-3:9092, properties.group.id flink-sales-stats, scan.startup.mode latest-offset, format debezium-json ); -- 1分钟滚动窗口统计每个SKU在每个门店的销量 CREATE TABLE sku_sales_1min ( window_start TIMESTAMP(3), shop_id STRING, sku_id STRING, total_qty INT, total_amount DECIMAL(12, 2), PRIMARY KEY (window_start, shop_id, sku_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:clickhouse://clickhouse:8123/retail, table-name sku_sales_1min, sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 5s ); INSERT INTO sku_sales_1min SELECT TUMBLE_START(event_time, INTERVAL 1 MINUTE) AS window_start, shop_id, sku_id, SUM(quantity) AS total_qty, SUM(sale_amount) AS total_amount FROM pos_sales GROUP BY TUMBLE(event_time, INTERVAL 1 MINUTE), shop_id, sku_id;这段SQL的实现逻辑很直白但真正跑起来之后发现状态管理是个大问题。每个SKU每个门店每分钟一行结果如果全是热销品还好但我们平台上总共30万个SKU分布在上百家门店时间窗口一开Flink的状态后端压力就会飙升。我们后来做了两步优化一是引入State TTL默认48小时自动清理过期状态二是开启Changelog机制减少CheckPoint的大状态序列化开销。这里要特别提醒Flink的CheckPoint间隔最好不要设置太短1分钟一次就够了否则频繁做全量快照会导致整个链路出现周期性延迟。3.4 数据一致性保障Exactly-Once语义在库存扣减中的实践实时计算里最难的问题不是快而是准。库存扣减这个动作对一致性要求极高——不能多扣会导致超卖也不能少扣会导致库存虚高。Kafka默认的at-least-once语义配合Flink的CheckPoint机制理论上能实现exactly-once但前提是下游Sink也必须支持幂等写入或者事务性写入。ClickHouse本身不支持传统意义上的事务所以我们采用了幂等写入 去重表引擎的方案。具体做法是在ClickHouse里使用ReplacingMergeTree引擎以(window_start, shop_id, sku_id)作为去重键重复写入的数据会被覆盖。同时在下游应用层维护一个processed_offset表记录每个分区已处理到的offset。这样即使Flink任务重启并发生回溯重复写入的数据也不会造成重复计数。还有一点必须提的是数据回放与对账机制。每天凌晨2点Spark批处理任务会重新计算前一天所有SKU的库存变更明细与实时链路产出的结果做全量对账。一旦发现差异比如某个SKU的期末库存对不上就会触发告警并自动修正ClickHouse中的指标数据。这个对账机制让业务方对我们的实时数据从不敢用到放心用信任是一点一点建立起来的。4. 库存与供应链优化模型落地算得出来还要送得到4.1 需求预测与安全库存计算不是拍脑袋拍出来的是算出来的数据链路打通之后真正产生业务价值的是优化模型。我们的核心模型是动态安全库存计算它直接决定了每个SKU在每家门店应该放多少货。传统的安全库存公式是SS Z * σ * √(LT)其中Z是服务水平对应的安全系数比如95%服务水平对应Z1.65σ是需求的标准差LT是补货提前期。但这个公式在零售场景里有一个致命缺陷它假设需求服从正态分布但零售销量尤其是促销期的销量往往呈现长尾分布和明显的日周期性。所以我们引入了分位数回归的方法直接预测需求的P95分位数而不是用平均值和标准差去推导。这样算出来的安全库存对促销期的需求尖峰更敏感有效降低了缺货率。预测方面我们用了一个比较务实的方案Prophet基础预测 实时销量修正。Prophet对周期性和节假日效应的建模比较友好在米兰这种节假日较多的城市很适用。但Prophet本身不支持在线学习所以我们在它上面套了一层修正机制实时计算模块统计当天截至当前的销量与Prophet预测的全天销量对比如果偏差超过阈值就用一个修正系数去调整当前的安全库存和补货建议。这套机制上线后米兰三家旗舰店的缺货率从平均8.2%降到了3.5%效果相当明显。4.2 补货建议的决策逻辑从拍脑袋到有据可依补货建议不是简单地库存低于安全线就补货那会导致很多无效补货和资金占用。我们的决策引擎集成了一系列规则和优先级SKU分级先按销售额和毛利贡献度将SKU分为A/B/C类。A类SKU贡献80%销售额的那部分补货频率高、补货量大C类SKU则采用最小补货量策略避免积压。仓配约束每个门店有固定的仓库容量上限以库位数量计。补货建议必须考虑当前门店的剩余库容不能超出物理上限。供应商交付周期不同供应商的交期差异很大有的2天能到有的要10天。补货建议必须根据交期动态调整订货点交期长的SKU要留更高的安全库存。季节与促销因子系统维护了一个促销日历大促前N天会提前提高安全库存阈值防止促销开始后库存被打穿。补货建议生成后并不是直接下采购单而是推送给采购员做确认。这里有一段系统工作流的设计经验人机协同比全自动更靠谱。采购员看到的是系统建议决策理由比如该SKU过去7天日均销量增长40%当前库存仅够支撑2天由人来做最终判断。上线三个月后采购员对系统建议的采纳率从最开始的60%提升到了90%原因是系统建议越来越贴合实际情况逐步建立了信任。4.3 跨店调拨与订单分配如何把消息变成动作再变成收益跨店调拨是高库存周转的核心手段之一。A店缺货、B店积压系统需要在15分钟内识别这种不平衡给出调拨建议。调拨的目标函数是最大化整个门店网络的总利润要平衡的是调拨成本物流费用、损耗风险和缺货损失失去销售机会。实现上我们使用了一个简化的混合整数规划模型。决策变量是每个SKU从B店调拨到A店的数量约束条件是B店的库存下限、A店的需求量上限、物流车辆的载重限制。目标是最小化缺货损失 调拨成本之和。这个模型每天跑两次上午10点和下午4点每次求解时间控制在3分钟以内求解器用的是开源SCIP。对门店数量大、SKU多的场景模型会拆分成按品类、按区域的子问题并行求解避免求解时间过长。调拨建议生成后会被发送到运营人员的审批工作台确认后自动生成物流任务下发到TMS运输管理系统。这条链路的关键在于任务的闭环反馈——不只是给建议还要跟踪建议是否被执行、执行后是否真正解决了缺货问题。我们为每个调拨动作生成了唯一的transfer_id在物流签收后回写库存系统完成整个闭环。5. 压测暴露的问题和上线初期的五项关键系统调整5.1 压测方案设计与第一轮压测暴露的三个问题我们设计压测方案的时候没有按照技术团队自己的想象去构造数据而是找历史数据团队拿了去年黑色星期五和圣诞大促的实际流量曲线按1:1的比例回放再在此基础上叠加1.5倍的峰值系数。这样压测出来的结果才真的有参考意义。第一轮压测结果不算理想暴露了三个问题。第一个是Kafka消费者rebalance过于频繁导致Flink作业频繁进入恢复状态。排查后定位到是消费者心跳超时时间设置太短session.timeout.ms默认是45秒但我们消费端频繁GC导致心跳暂停超过阈值。解决方式是调大session.timeout.ms到120秒同时优化Flink的内存配置减少Full GC频率。第二个问题出在ClickHouse的写入毛刺上。高峰期写入速率超过每秒2万行的时候ClickHouse的merge线程忙不过来出现写入排队。我们用Buffer表引擎做了一层缓冲让数据先写入Buffer表再异步刷到主表同时增加了background_pool_size和background_merge_pool_size。调整之后写入毛刺消失了。第三个问题挺隐蔽Redis热键集中导致集群节点CPU倾斜。Top100的热销SKU数据全在一个节点上每秒几万次的读请求把这个节点的CPU打到接近100%。后来我们用一致性哈希加虚拟节点的方式做分片同时把热SKU的数据做了多副本缓存才算把问题解决。5.2 上线初期监控告警策略的演进与事故复盘系统上线第一周我们就经历了一次真实的库存数据延迟事故。周五晚上8点是米兰门店的销售高峰期线上订单量突然翻了4倍一个折扣活动提前开始结果Flink作业反压严重从数据进入到指标产出延迟从正常的5秒左右拉长到了6分钟。库存大屏上的数据和真实库存对不上门店店长立刻就发现了问题。事后复盘根因有两点一是Flink作业的并行度没有针对这种突刺流量做弹性扩容二是监控告警的阈值设得太宽了延迟超过5分钟才告警。我们做了三个改进第一把Flink作业的auto-scaling打开基于Kafka消费延迟和TaskManager的繁忙度指标自动调整并行度第二告警阈值收紧到延迟超过30秒就触发P1告警并增加电话通知第三增加了一个独立的看门狗服务每秒探测一次数据链路末端的最新事件时间watermark如果发现watermark落后当前时间超过60秒就自动触发链路降级——把非核心计算任务暂停优先保障库存扣减和销量统计这两个核心任务。这三个改进落地之后之后即便是黑五当天系统最大延迟也没超过20秒核心数据链路始终保持健康。5.3 成本调优实践高并发不等于高成本最后说说成本问题。高并发数据链路如果设计不好云资源费用会非常可观。项目上线三个月后我们做了一轮全面的成本Review主要优化了三个方面。第一是冷热数据分层存储。ClickHouse里超过30天的明细数据不再保留热存储而是定期迁移到对象存储S3查询时通过Disk类型的表引擎远程读取。这个改动直接让ClickHouse的存储成本下降了60%左右。第二是Flink资源按峰谷调整。白天10点到晚上10点的流量高峰维持32个TaskManager夜间低峰时段缩减到8个。配合Kubernetes的HPA自动伸缩计算资源成本大约节省了35%。第三是Kafka的日志留存策略调整。原来所有topic默认保留7天数据实际上大部分数据在实时计算消费完之后就没用了。我们按topic的重要程度分别设置了不同的retention核心业务topic保留72小时用于对账和回溯日志类topic只保留24小时。存储成本又降了一截。这些优化做完之后整个实时数据链路每月的云资源成本降幅接近40%而系统性能和稳定性指标没有下滑。对于任何一个需要在预算约束下长期运营的数据平台来说省钱还能稳住性能这件事本身就是核心竞争力。6. 踩过的坑和留给后来人的实操建议这个项目做下来有几件事如果让我重新选或者跟同行分享我会单独写出来。首先是不要一上来就追求全链路实时零售业务里真正要求秒级响应的场景很少。我们在做需求梳理时发现大约只有20%的数据需要真正的实时处理剩下的80%其实走准实时分钟级就够了。如果一开始就把所有数据都往Kafka Flink这条链路上赶不仅资源浪费运维复杂度也会倍增。先分清楚实时、准实时和批处理的边界再动手设计架构这是最值得花时间的一步。其次是数据契约比数据平台本身更值得投入。我们在项目过程中多次遇到上游系统调整了字段含义下游却毫不知情导致计算结果悄悄出错。后来建立了严格的数据契约评审机制每次上游变更必须通过契约评审涉及核心指标字段的变更必须做灰度验证。这条流程虽然官僚了一点但对保障数据质量起到了决定性作用。最后是关于团队协作方式的一点体会。这个项目需要数据工程师、算法工程师、供应链业务专家和后端开发紧密配合但每个人关注的东西完全不同。业务专家关心的是建议是否合理算法工程师关心的是模型准确率如何数据工程师关心的是链路是否稳定。如果大家没有共同的目标和统一的衡量标准很容易互相甩锅。我们每周固定两次站会每次只围绕三张报表看缺货率、库存周转天数、数据链路延迟。目标简单分歧就少效率自然就高了。系统上线到现在我最大的感受是高并发数据平台不是指标的堆砌不是Kafka吞吐量多高、Flink并行度多大、ClickHouse查询多快而是这些技术最终有没有让门店的货更准、让仓库的货转得更快、让采购的人少拍脑袋。技术指标只是手段业务结果才是目的。如果正在读这篇文章的你也在做类似的项目代码问题、架构问题都好解决最难的是让整个业务链条里的人真正信任并且用起来这套系统。信任不是一个晚上能建立的但只要数据足够准、链路足够稳、建议足够实用团队会一点点从怀疑变成依赖。这也是这套系统做到现在最让我觉得有价值的地方。