1.2 摄取、存储与保留:把数据管道做成可重放的系统
契约说明了一条记录代表什么。接下来,档案管理员把你带到预言厅地下:队列会重试,批任务会部分失败,历史分区会被回补。这里的数据生命周期不是七个静态名词,而是一组必须能够重放、观察和删除的运行流程。
本课目标
- 根据延迟、吞吐和一致性选择批处理或流处理;
- 设计 raw、validated 与 serving 层之间的边界;
- 用 checkpoint、幂等和 reconciliation 处理失败;
- 建立保留、归档与删除的可验证流程。
1. Batch 与 streaming 是数据边界选择
批处理面对有界数据集:某天文件、某次快照、一个日期区间。流处理面对持续到来的无界事件,系统用窗口和 watermark 形成暂时可计算的边界。
| 维度 | Batch | Streaming |
|---|---|---|
| 结果延迟 | 分钟到天 | 毫秒到分钟 |
| 输入边界 | 明确结束 | 持续到达 |
| 失败恢复 | 重跑批次 | checkpoint + replay |
| 迟到数据 | 下一批回补 | watermark/allowed lateness |
| 运营复杂度 | 通常较低 | 状态、顺序与背压更复杂 |
同一系统可以“流式摄取、批量分析”,也可以用微批实现近实时。业务需要五分钟内告警,不意味着全部历史分析都必须改成流处理。
2. 摄取先建立批次身份与清单
读取每日 CSV 时,至少记录:
batch_id
source_uri
source_version / etag
expected_size
checksum
observed_rows
schema_version
started_at / completed_at
statusfrom pathlib import Path
import hashlib
import pandas as pd
path = Path("missions_snapshot_2026-08-02.csv")
digest = hashlib.sha256(path.read_bytes()).hexdigest()
frame = pd.read_csv(path)
manifest = {
"source": str(path),
"sha256": digest,
"rows": len(frame),
"columns": list(frame.columns),
}
print(manifest)示例为说明 manifest,会把整个文件读入内存计算 hash;大文件应分块计算。df.info() 自己打印信息并返回 None,不要写 print(df.info()) 期待得到结构化报告。
3. 端到端核对比网络丢包口号更有效
TCP 或消息中间件可能处理传输级重试,但业务数据仍会因生产事务、序列化、过滤、消费者提交时机和目标写入失败而丢失或重复。
可核对:
- 源侧记录数与目标接受/拒绝数;
- event ID 去重前后数量;
- 各分区最小/最大时间;
- 金额、计数等守恒汇总;
- dead-letter 数量与原因;
- source offset 与 checkpoint。
简单行数相等不是充分证明:一条丢失加一条重复仍可能总数相同。稳定 ID、校验和与业务汇总要组合使用。
4. 分层存储隔离责任
一种常见但非唯一的分层:
raw 保留接收到的原始表示与 provenance
validated 通过 schema/基础质量检查,坏记录隔离
curated 语义统一、去重、业务规则明确
serving 针对报表、API 或模型组织每层应有明确输入、输出和重建方式。不要把“bronze/silver/gold”当作固定真理;名称不重要,可重放边界与质量承诺才重要。
存储选择由访问模式、更新语义、规模、延迟和成本共同决定。列式文件适合扫描分析,但小文件过多会拖累元数据和调度;行式数据库适合点查与事务,却不意味着不能聚合。不要仅凭“按天聚合”就断言必须换某种存储。
5. 处理步骤必须幂等
任务在写到一半时失败,重跑不能把结果重复累加。常见策略:
- 按分区写临时位置,成功后原子发布;
- 以业务键执行 upsert;
- 输出路径包含确定性 run/version;
- checkpoint 与外部副作用协调;
- 记录输入版本到输出版本的 lineage。
“先删目标分区再重写”虽然简单,却会产生暂时空窗,并在删除后失败时丢失旧结果。更安全的是写新版本、校验、再切换元数据指针。
6. Backfill 与日常运行使用同一逻辑
修正历史 bug 时需要回填。若回填另写临时脚本,规则很快漂移。管道应参数化:
process(start_time, end_time, input_version, code_version)回填前评估:
- 会重写哪些分区;
- 下游是否自动触发;
- 当前数据是否混用新旧规则;
- 资源是否挤占在线任务;
- 如何比较旧版与新版结果;
- 如何回滚发布指针。
7. 可观测性关注数据状态
除 CPU、内存和任务成功率外,还要监控:
- freshness:最新完整数据距现在多久;
- volume:记录数、字节数与历史基线;
- validity:schema、范围、枚举违规;
- completeness:关键字段与分区是否齐全;
- uniqueness:业务键重复;
- distribution:分布是否发生异常漂移;
- reconciliation:源与目标是否守恒。
阈值要考虑工作日、季节性和业务发布。仅以“比昨天少 10%”告警,会在周末稳定制造噪声。
8. 保留和归档由义务驱动
保留期应综合:
- 产品与分析需求;
- 法律保存义务与删除权;
- 调试和审计窗口;
- 重放所需历史;
- 存储与恢复成本。
归档前要验证可读性:对象存在不代表几年后仍有 schema、密钥和软件能够解释。定期做恢复演练,并保存格式、schema 与校验信息。
压缩比例取决于数据分布和编码,不能承诺统一“减少 80%”。应对样本测试 gzip、zstd、Parquet 编码等方案,比较成本与读取模式。
9. 删除必须覆盖派生与备份
删除地图包括:
原始对象 → 清洗表 → 聚合表 → 特征 → 搜索索引 → 缓存 → 备份删除文件名不等于物理介质立即不可恢复。现代 SSD、对象存储和托管数据库的擦除由底层实现决定,反复覆写未必可靠或可用。常见方法是:
- 逻辑删除后由存储生命周期清除版本;
- crypto erasure:销毁独立加密密钥;
- 使用供应商保证的介质清退与擦除流程;
- 对备份设置不可变但有限的过期窗口;
- 保存不可含敏感内容的删除审计证据。
法律上需要删除什么、何时完成取决于适用规则和角色,应由合规与安全团队确认,不能由教程给出统一期限。
常见误区
- 任务显示成功就代表数据完整:成功状态不覆盖语义核对。
- 行数相等就没有丢失:重复可以抵消缺失。
- 原始层永远保留最安全:会扩大敏感数据与合规风险。
rm或反复覆写适用于所有存储:介质、版本和托管层行为不同。
练习
- 为每日快照设计 batch manifest 和重跑幂等策略。
- 构造“行数相等但一丢一重”的核对反例。
- 为一个历史回填写发布、验证与回滚步骤。
- 画出用户删除请求需要覆盖的所有派生数据节点。
小结
可运营的数据管道要能说明输入版本、处理状态、输出版本和失败恢复方式。批流选择决定边界,幂等与重放处理故障,质量监控判断结果是否可信,保留与删除则控制数据在何时退出系统。
下一章进入清洗工作台,但不会从“空值填中位数”开始,而是先识别缺失机制、类型歧义和每条修复规则正在做出的假设。