搞数据工程这些年,被问得最多的问题就是:能不能推荐一个适合练手的真实项目?网上那些例程,不是对着几张假表抽来抽去,就是动不动就上集群,学完之后连业务方要什么都不知道。我这个无人售货机零售数据项目,算是一个折中的选择:数据量不大,但五脏俱全,订单流水、设备状态、库存快照、补货记录全都有,支付失败、退款、设备掉线、库存对不上这些零售行业真正的脏数据,一个都不少。拿它来练ETL,基本能把数据接入、清洗、建模、调度的完整流程跑通。这篇就是我带这个项目时的完整实践记录,适合正在学数据仓库、数据集成,或者想从SQL开发往数据工程方向走的同学参考。
1. 项目设计拆解:无人售货机零售数据到底长什么样
1.1 业务背景:售货机背后有几个系统
无人售货机从外表看就是一个铁柜子,但背后牵扯的系统比想象中多。用户扫码下单,机器出货,资金清算,补货员定期补货,运维人员远程看状态,这些都是不同系统在配合。做ETL项目如果只盯着订单表,后面分析时缺设备维度、缺库存信息,还得返工。所以我的建议是第一步先把业务链路拆清楚。
一台售货机至少涉及四类系统,分别是交易系统、设备系统、库存系统和供应链系统。交易系统负责订单和支付,设备系统负责心跳上报和设备状态,库存系统维护货道实时库存,供应链系统记录补货和维修工单。它们之间的数据流转关系是:用户在设备上扫码支付,交易系统生成订单并通知设备出货;设备出货后上报库存变化给库存系统;补货员补货时在供应链系统创建补货单,同时更新库存;设备系统则一直在上报心跳,确认机器在线。
做ETL之前,把这张业务关系图画在脑子里非常有必要。你后面建表、写清洗逻辑、做数据校验,全都要靠这份业务理解兜底。如果只关心订单表的主键和金额,那充其量是个SQL取数员,不是在做数据工程。
1.2 为什么拿售货机练ETL
有人会问,电商订单数据量大,练起来不是更带感吗?但个人练手和工业生产环境是两回事。电商订单动辄一天上千万行,没有集群根本跑不动,一台普通电脑压根撑不住;而售货机的数据量级刚刚好。单台机器一天通常几十到几百单,一个城市几十台机器,一天大概几十万行,这个量级用一台笔记本、装个MySQL加Python环境就能完整跑通。
更重要的是,数据量虽小,问题种类一点不缺。售货机数据里有跨时区时间、设备离线补传、订单取消退款、商品改价历史、货道库存出现负数、同一台机器接入了微信支付宝银联等多个支付渠道。这些恰恰是真实企业数据环境里最常见的问题。ETL的核心能力不是会写几个SQL,而是面对这些千奇百怪的数据还能保证结果准确、稳定。拿售货机数据练手,成本低,见效快,练的就是这个基本功。
1.3 项目目标与交付物
这个项目的目标是做一套完整的离线数仓链路,覆盖从业务库到分析表的全过程。具体交付物有这样几个:
- 统一的ODS原始数据层,保留业务库原样数据;
- 规范的DWD明细层,完成清洗、去重、标准化;
- 商品、设备、日期等维度表,支撑多维分析;
- 订单事实表、库存快照表、补货事实表;
- 按天调度的ETL任务,可重跑、可监控;
- 一套销售、库存、设备运营分析指标。
我建议做这个项目时不要急着写代码,先把交付物列清楚。我在带项目时发现,凡是先花时间做设计的,后面都会顺利很多;一上来就写抽取脚本的,基本都会在中途推翻重来。
2. 数据源梳理与采集通道设计
2.1 五个关键数据实体和字段说明
整个售货机业务链路里,最核心的数据实体有五类,分别是订单流水、商品主数据、设备档案、库存快照、补货记录。我把它们的来源和关键字段整理成了一个表:
| 数据实体 | 来源系统 | 典型字段 | 更新频率 | 相对量级 |
|---|---|---|---|---|
| 订单流水 | 交易系统 | order_id, device_id, product_id, channel_no, pay_type, amount, status, order_time | 分钟级,持续产生 | 最大 |
| 商品主数据 | 商品中心 | product_id, product_name, category_id, price | 低频,偶尔修改 | 很小 |
| 设备档案 | 设备平台 | device_id, location, model, owner_id, status | 低频 | 很小 |
| 库存快照 | 库存系统/设备上报 | device_id, channel_no, stock_qty, snapshot_time | 每小时或每次变动 | 中等 |
| 补货记录 | 供应链系统 | replenish_id, device_id, operator, replenish_time, replenish_qty | 每日少量新增 | 小 |
这五类数据的更新频率差异很大。订单和库存是持续产生的高频数据,商品和设备是基本不变的维表。做采集设计时必须区别对待:不能把所有数据都按天全量拉,也不能把静态维表做成高频增量。订单表适合增量抽取,维表适合定期全量刷新或拉链处理,库存快照则根据设备上报频率决定采集周期。
2.2 三种抽取通道怎么选
无人售货机后端在实际项目中通常有三种数据输出方式,我建议把三种方式都过一遍,因为企业里真的会遇到“部分数据只能靠文件”的情况。
第一种是业务库直连抽取。数仓平台直接连业务MySQL或者PostgreSQL,通过SQL按条件查询数据。这种方式实现最简单,成本最低,适合离线数仓的日常调度。缺点是会占用业务库的连接和IO资源,对高并发业务库有影响,所以查询条件必须严格过滤,尽量走索引,不要全表扫描。
第二种是消息队列订阅。交易系统在订单创建、支付成功、退款完成等关键节点,把事件写入Kafka这类消息队列,数仓侧通过消费者订阅这些消息,解析后落库。这种方式适合做近实时链路,也能减轻业务库压力,但需要处理消息乱序、重复消费、积压等问题。
第三种是文件接口。部分老旧的售货机设备或第三方支付渠道,会定期生成Excel或CSV对账单,通过SFTP或OSS上传。这类数据通常用于和交易系统的订单做对账、补漏。做数据接入时,不能因为它是文件就觉得低级,对账场景里它反而是最重要的数据源。
2.3 调度频率与增量策略
不同表的数据刷新策略,我建议按这个原则来定:事实表按增量走,维表按全量或拉链走,小表全量覆盖,大表增量追加。
订单表用增量抽取,记录每一批次的最大处理水位,下次只捞水位之后的数据。库存快照表同样是增量追加,但因为存在设备离线补传的情况,时间字段要做特殊处理。商品和设备维表采用每日全量快照覆盖,在数仓里保留当天的全量版本即可。补货记录一天没几条,全量拉取也没有压力,按天覆盖就行。
调度频率上,日批任务是底线。订单和库存这两个表建议做到小时级或者半小时级调度,这样运营第二天看数据时能看到前一天的全量数据,当天白天也能看到实时进展。我常用的调度组合是:每半小时拉一次订单,每小时拉一次库存快照,每天夜里拉一次维表和补货记录。这个频率对一台普通服务器完全没有压力,但已经能满足大多数售货机零售业务的分析需求了。
3. 转换与建模:从业务表变成分析宽表
3.1 数据清洗五个高频场景
数据抽取到数仓后,还不能直接用于分析,必须做清洗和标准化。我在售货机项目里遇到的清洗场景主要集中在五个方面。
第一个是时间字段格式和时区不统一。设备端时间、服务端时间、支付回调时间经常不在一个时区,有的带UTC后缀,有的带+08:00,有的干脆是秒级Unix时间戳。建议统一转成东八区时间,字符串格式统一为YYYY-MM-DD HH:mm:ss。这一步要放在抽取之后立刻做,后面所有表都用这个标准时间。
第二个是金额字段出现字符串、负数或科学计数法。售货机的订单金额理论上不会出现负数,但退款场景会产生冲正记录。所以设计订单表时,要把订单金额、支付金额、退款金额分开存储,不能混在一个字段里。
第三个是支付状态枚举值混乱。有的系统用0、1、2表示状态,有的用SUCCESS、FAILED、CANCELLED,还有的直接用中文“成功”“失败”“已退款”。如果不在DWD层统一映射,下游做统计时根本没法用。我通常会把状态映射表单独建一张维表,代码里用CASE WHEN或者字典映射来转换。
第四个是设备ID和商品ID在不同系统不一致。同一台机器,设备平台叫DEV001,交易系统却叫10001,如果不做ID归一化,两张表JOIN不上,后面所有分析都会出错。这个问题的解决办法是建立一张设备映射表,把业务侧ID和数仓侧代理键对应起来。
第五个是重复订单。网络超时重试、消息重复消费都会造成同一订单出现多条记录。抽取阶段必须做幂等去重,去重键用业务订单号,不能使用数据库自增ID。
3.2 星型模型和核心表设计
数仓设计我推荐走星型模型。售货机业务复杂度不高,完全不需要搞宽表大而全,订单一张事实表,周边挂商品、设备、日期三张维表就足够了。库存快照和补货记录是独立的业务过程,单独建表,不用硬塞进订单宽表。
订单事实表的设计我通常采用这样一组字段:订单ID、设备ID、商品ID、货道号、支付类型、订单金额、优惠金额、支付金额、订单状态、下单时间、支付时间、退款时间、分区日期。其中货道号这个字段比较容易被忽略,但它对后续分析很有价值,比如可以分析哪些货道动销差、哪些商品摆放位置影响销量。支付类型要区分微信、支付宝、现金、会员卡,后面做渠道分析时用得上。优惠金额和支付金额分开存,是为了对账时能还原每一笔订单的真实流水。
日期维表不要省。日期维表里应该有日期、年、月、周、星期几、是否节假日等字段。有了它,随便写一条GROUP BY date JOIN维度表的SQL就能按周、按月做趋势分析,非常方便。我见过不少项目省了日期维表,结果后面每个查询都要自己写日期格式化,重复劳动不说,口径还不统一。
3.3 商品维度的拉链更新
商品的价格和名称是会变的,饮料涨个价、包装改个版,这在零售行业太常见了。如果维表直接UPDATE主键,历史订单关联到的商品信息就会被篡改,后续做销售分析时口径就乱了。所以商品维表要做成拉链表,保留每条记录的有效时间区间。
最简单的拉链实现方式是:每天全量拉取商品主数据,和昨天的快照做对比。如果主键相同且属性没有变化,不处理;如果属性有变化,就把旧记录的失效时间置为昨天,新记录从今天开始生效;如果是新增商品,直接插入一条新记录。这样一个商品可以有多个版本,但在任意时间点上只有一条记录处于有效状态。做历史销售分析时,按订单日期JOIN商品维表,就能回算出当时的价格、分类和商品名。
4. ETL全流程实操记录
4.1 订单增量抽取编写与参数控制
订单数据在业务库中持续产生,每次全量抽取不现实,增量抽取是唯一正确的做法。增量抽取最稳妥的方式,不是每次全量对比,而是记录上一批次处理到的最大update_time水位线。
以MySQL为例,订单表通常会有update_time字段,但这个字段必须保证和订单状态变更联动,也就是说订单创建、支付成功、退款完成时都必须更新update_time。抽取SQL可以写成这样:
SELECT order_id, device_id, product_id, channel_no, pay_type, amount, discount_amount, payment_amount, status, order_time, pay_time, refund_time, update_time FROM business_order WHERE update_time > '{last_max_update_time}' AND update_time <= '{current_batch_time}' ORDER BY update_time;这里有两个关键细节。第一,抽取边界必须是半开区间,左边严格大于上次水位,右边小于等于当前批次时间,两边如果都是大于等于,两个批次之间就会重复。第二,抽取完成后要把当前批次的最大update_time记录到调度元数据表里,这个值就是下一次任务的起始水位。
如果业务库的update_time字段更新不及时,导致增量漏数据,那就需要加一道兜底逻辑:额外对当天全天的订单做一次完整性校验,比较业务库当天订单数和数仓当天记录数,不一致时触发全量重抽。
4.2 明细层装载与幂等去重
数据从业务库抽取出来之后,要先落到ODS层,也就是原始数据层。ODS层的作用是保留证据,字段和业务库保持一致,不做任何清洗。为什么要保留这一层?因为之后如果发现DWD层的清洗逻辑写错了,还能从ODS原始数据重新计算,不然数据一旦被覆盖就无法追查了。
ODS层落完之后,才是DWD层的清洗和标准化。在写入DWD订单表之前,要做幂等去重。我的做法是先把ODS当天增量数据加载到一个临时表,然后按业务订单号去重,保留更新时间最新的一条,再INSERT到DWD明细表。
去重的SQL可以这样写:
INSERT INTO dwd_order_fact SELECT t.order_id, t.device_id, t.product_id, t.channel_no, t.pay_type, t.amount, t.discount_amount, t.payment_amount, t.status, t.order_time, t.pay_time, t.refund_time, t.dt FROM ( SELECT *, ROW_NUMBER() OVER ( PARTITION BY order_id ORDER BY update_time DESC ) AS rn FROM ods_order_data WHERE dt = '{batch_date}' ) t WHERE t.rn = 1;这个去重逻辑看似简单,但实际价值很大。我在项目里遇到过不止一次因为消息重复推送导致同一订单出现两条记录的情况,如果没有这层防护,订单金额会被翻倍统计,业务报表直接出错。
4.3 库存快照与补货事件的加工
库存快照是售货机数据分析里很容易被忽略但非常关键的数据。它记录的是每个货道在某个时间点的库存数量。售货机的库存和电商不一样,电商库存是账面库存,售货机库存是物理库存,两者经常不一致。
为什么会有差异?最典型的原因是货道卡货。系统判断出货成功并且扣减了库存,但商品实际卡在货道里没有掉落到取货口,用户没拿到货,货道里却有货,账面和实物就对不上了。还有一种情况是补货员补货时,系统库存没有同步增加。这些都会导致库存快照和理论库存产生偏差。
所以ETL里不能只是把库存快照存起来展示,还要做联动校验。我的做法是把销售记录和补货记录合并计算理论库存,再和设备上报的库存快照做对比:
- 期初库存 + 补货数量 - 销售数量 = 理论库存;
- 理论库存 与 设备上报库存 的差异超出阈值时,生成异常记录;
- 异常记录推送给运营,安排现场盘点。
这套逻辑跑起来以后,缺货预警、补货计划就都有数据支撑了。
4.4 调度串起来的完整效果
把上面所有脚本串成一个可调度的流水线,理想状态大概是这样的:
- 每30分钟调度一次订单增量抽取任务,抽完落到ODS,再清洗进DWD;
- 每60分钟调度一次库存快照抽取任务,同步进行理论库存校验;
- 每天凌晨2点调度商品维表、设备维表的全量拉取和拉链更新;
- 每天凌晨3点调度补货记录抽取和汇总指标计算;
- 每天凌晨4点输出经营日报,包括销售额、订单量、库存预警、设备异常等。
我在实际项目中用Airflow来做调度,每个任务之间设置依赖关系。订单清洗必须在上一次订单抽取成功后才能运行,库存校验必须等库存快照和订单明细都就绪后才能运行。任务失败时要能自动重试并告警,我通常设置重试2次,每次间隔5分钟。
5. 踩坑记录与问题排查实录
5.1 订单时间全部偏移8小时
这个坑我印象太深了。第一次跑数,发现当日订单比实际少了将近三分之一,排查了很久,最后定位到设备端上报时间存的是UTC时区,但数据库连接时区设置不对,写入时直接把UTC当作东八区存了,导致所有订单时间都晚了8个小时。凌晨的订单跑到了当天上午,上午的订单跑到了下午。
这类问题的修复不能指望业务系统自己改,必须在ETL层做兜底转换。用Python处理时可以这样写:
import pandas as pd df['order_time'] = pd.to_datetime(df['order_time'], utc=True) df['order_time'] = df['order_time'].dt.tz_convert('Asia/Shanghai') df['order_time'] = df['order_time'].dt.strftime('%Y-%m-%d %H:%M:%S')从那以后,我养成了一个习惯:所有时间字段在数据入库前统一打印一批样例,肉眼确认时区偏移是否正常。
5.2 重复消费导致重复订单
用消息队列接入订单数据时,如果消费者处理完数据但没有提交偏移量,网络一抖动,下次就会重新消费同一批消息,导致订单表出现重复记录。这个问题一开始没被发现,直到月度对账时发现流水总额比支付渠道多了十几万才排查出来。
从此以后,我要求所有事实表加载任务必须带上去重逻辑。不管抽取源头是什么,加载时都要按业务主键做一次幂等校验,优先保证同一订单只保留一条记录。这是做数据接入的底线,不能省。
5.3 货道库存出现负数
库存显示-3,但货道里明明有货,这个现象一度让运营非常困惑。排查后发现是设备上报库存快照和交易系统扣减库存的时序不一致:设备先发生了商品掉落,但上报库存的动作延迟了,随后又产生了交易扣减,导致快照计算出现负数。
处理上,我建议设置库存下限为0,并把异常快照单独打标,不去直接覆盖库存字段,保留原始快照和修正值两条记录。这样后续盘点时还能追溯是哪台设备、哪个货道、什么时候出现的异常。
5.4 商品名称对不齐
售货机商品名称经常出现“可口可乐330ml”“可乐(罐装)”“可口可乐迷你罐”这种写法,理论上是同一商品,但因为来源名称不一致,直接按商品名分组统计会得到一堆冗余分类,后续月度趋势分析根本没法做。
处理思路是建立一套商品名称映射规则,核心是用商品ID做关联,商品名只用于展示,不做聚合Key。如果业务数据里没有商品ID,就只能靠正则匹配和历史沉淀的映射表来归并。这个工作很枯燥,但必须做,否则数据质量永远在及格线以下。
5.5 问题排查速查表
我把这个项目里踩过的典型问题整理成了一个速查表,方便后面复用:
| 问题现象 | 可能原因 | 排查思路 | 解决办法 |
|---|---|---|---|
| 当日订单量偏少 | 时区偏移、抽取水位线设置错误 | 对比业务库当天总数和数仓总数 | 统一时间转换,检查增量水位 |
| 订单金额翻倍 | 重复消费或重复抽取 | 按订单号统计重复条数 | 加载前去重,按业务主键幂等 |
| 库存为负 | 库存上报与扣减时序不一致 | 核对设备上报时间和交易时间 | 下限置0,异常快照打标人工复核 |
| 商品分类过多 | 商品名称不统一 | 统计同一商品不同名称 | 建立商品映射表,按商品ID关联 |
| 任务偶发失败 | 业务库连接数超限或网络抖动 | 查看任务日志和数据库连接配置 | 设置自动重试和告警,抽取加索引条件 |
| 设备离线无数据 | 机器断电或网络断开 | 查看设备心跳表 | 用补货记录和销售记录补重建库存 |
6. 后续扩展方向:从离线批处理到实时数据服务
6.1 实时化改造思路
这个售货机项目做完离线链路之后,后续可以往实时方向扩展。订单表和库存快照更新频率很高,天然适合尝试实时数仓,比如用Flink或Spark Structured Streaming把订单数据流式接入,实时统计当前设备销售额排名、各货道动销情况、缺货预警。
做实时化改造的时候,第3节做的ID归一化、时间统一、维度映射逻辑可以直接复用,不用再重新踩一遍坑。唯一要额外处理的是消息乱序:实时流里订单支付成功的时间可能比订单创建时间先到,所以要基于事件时间而不是处理时间来做窗口计算。
6.2 智能补货与异常预警
另一个很实用的扩展方向是预测补货。用历史订单数据和补货记录,可以按设备、按货道统计销售规律,再结合星期、天气、节假日等外部因素,预测每台设备未来三天的销量,输出补货建议清单。补得好不好直接影响运营成本,补太多货砸在机器里浪费周转空间,补太少又损失销售机会。
设备健康监控也值得做。把设备心跳上报的成功率、交易成功率合并到一个运营看板里,设备连续一段时间没有心跳或者交易成功率骤降时自动告警,能显著降低人工巡检成本。这些能力在离线数仓的基础上逐步叠加,项目就从一个教学案例逐步长成一个完整的数据产品了。
这个项目我带了好几批人,最大的感受是:ETL的难点从来不在工具本身,而在对业务数据的理解和对细节的把控。很多人一上来就写代码,结果建的表没法用,重来好多遍。如果你也打算拿这个项目练手,建议先花一天时间,把业务表和字段含义彻底摸清楚,再动工写脚本。最后分享一个伴随我多年的小习惯:每次任务跑完后,把当批次抽取行数、去重行数、异常行数打印成日志。多写这几行日志,排查问题的时候能省下成倍的时间,这个习惯我从售货机项目一直沿用到了后来的生产环境里。