S04:迟到、乱序、重复、缺失与 Schema Drift
数据时序与结构异常不是单纯的“脏数据”,而是会改变窗口聚合、标签、特征与决策含义的系统状态,需要显式定义处理语义。
内容类型:预习教材(不代表已完成)
日期:2026-11-26
阶段:P2 · AI Systems Engineering 90
总路线:Day 094 / 360
周次 / 节奏:W1 · 周四案例与连接
状态:教材已备;学习未完成
主题:event time、late data、deduplication、schema drift
一句话定义
数据时序与结构异常不是单纯的“脏数据”,而是会改变窗口聚合、标签、特征与决策含义的系统状态,需要显式定义处理语义。
学习目标
- 能区分 late arrival、out-of-order、duplicate、missing 与 schema drift。
- 能解释 watermark、去重键、缺失策略和兼容规则各自解决什么问题。
- 能分析这些异常怎样传播到训练、评测和线上决策。
- 能建立“问题→影响→处理→残余风险”表,而不是追求唯一正确答案。
核心知识
迟到表示事件到达时间晚于预期;乱序表示到达顺序与业务发生顺序不一致;重复表示同一事实被多次交付;缺失可能是字段未知、未采集或不适用;schema drift 则是结构或语义随时间变化。它们可以同时发生,例如旧版客户端晚到的一条重复事件既迟到又使用旧 schema。
窗口计算必须基于明确的时间轴。按 processing time 聚合简单,却会因网络和重放改变历史;按 event time 更接近业务事实,但要决定等待多久。watermark 是“系统认为更早事件大概率已经到齐”的进度假设,不是绝对保证。过晚事件可以丢弃、侧输出、更新历史结果或触发对账,选择取决于业务可逆性。
去重需要稳定的幂等键;相同 payload 不一定是重复交易,相同 event id 也可能因上游错误被复用。缺失值不能统一填 0,因为 0 是有含义的观察。schema drift 还包括类型未变但单位、枚举含义或采样逻辑改变的语义漂移。
机制与推导
对长度为 W 的事件时间窗口,统计量为:
[ A_t=\sum_i x_i,\mathbf{1}(t-W < event_time_i\le t) ]
若迟到事件在决策后进入,A_t 会从旧值变为新值。系统必须选择“历史真值可更新”还是“保留当时决策视图”。point-in-time 正确性要求训练样本在时刻 t 只能连接 available_time ≤ t 的特征,而不只是 event_time ≤ t;否则迟到但后来补齐的数据会穿越回过去。
对重复数据,简单计数偏差约等于重复事件贡献之和。去重状态又有保留窗口:保存太短会漏掉晚重复,保存太久会增加状态容量。因而正确性与资源存在取舍。
最小练习或观察步骤
- 准备 8 条合成事件,故意加入乱序、重复、缺字段、未来字段和一个极晚事件。
- 分别按文件顺序、processing time 与 event time 计算一个简单 24 小时总额。
- 设定一个假想 watermark,观察哪些事件被纳入、侧输出或需要修正。
- 为重复事件选择去重键,写出至少一个误杀和一个漏判反例。
- 填写“异常→训练影响→线上影响→候选处理→残余风险”表。
- 不实现流框架;重点解释状态与时序语义。
常见误区与边界
- 把 watermark 当成事件绝不会再到达的保证。
- 使用当前完整特征回填历史训练集,造成未来信息泄漏。
- 对所有缺失统一填 0,抹掉“未知”和“不适用”的区别。
- 只按 payload hash 去重,误删金额相同的合法事件。
- schema 能解析就认为兼容,忽略单位和业务含义变化。
系统场景连接
反欺诈系统在授权时只能看到当时已到达的交易与账户事实;次日补录信息可用于调查,却不应悄悄回填成模型当时已知的特征。链上数据还可能因区块确认与重组改变“已观察事实”。本日知识直接连接 feature store 的 point-in-time join、评测集构造和审计中的 decision-time snapshot。
自检问题
- late 与 out-of-order 的差别是什么?
- watermark 为什么是一项业务与容量共同决定的假设?
event_time ≤ t为什么仍不足以保证 point-in-time 正确?- schema 类型不变时还可能发生哪些语义 drift?
专业课程对齐
- 阅读 Feast 官方文档 中 point-in-time joins、historical retrieval 和 feature freshness,重点理解特征在事件时刻“可用”而非后来存在的要求。
- 阅读 OpenLineage 官方文档 的 dataset/run lineage,思考一次迟到数据回填怎样形成新的运行与输出版本,而不是静默覆盖历史。
- 阅读 Stanford CS329S 的 data distribution shift 与 production data 主题,将 schema drift 和数据分布变化区分开。
深入学习提示
先用纸面事件表手算,再看框架术语。重点画三条轴:事件何时发生、系统何时看到、决策何时作出。对每种处理策略都问“旧决策是否改写”“下游是否收到更正”“状态保留多久”。深入点不在于背 watermark 算法,而在于能说明正确性承诺、资源成本与无法消除的迟到风险。
学后填写区
- 实际构造的五类异常:
- 三种时间排序的结果差异:
- 选择的迟到处理语义:
- 去重的反例:
- 尚不确定的 point-in-time 边界: