跳到内容

1.2 摄取、存储与保留:把数据管道做成可重放的系统

契约说明了一条记录代表什么。接下来,档案管理员把你带到预言厅地下:队列会重试,批任务会部分失败,历史分区会被回补。这里的数据生命周期不是七个静态名词,而是一组必须能够重放、观察和删除的运行流程。

本课目标

  • 根据延迟、吞吐和一致性选择批处理或流处理;
  • 设计 raw、validated 与 serving 层之间的边界;
  • 用 checkpoint、幂等和 reconciliation 处理失败;
  • 建立保留、归档与删除的可验证流程。

1. Batch 与 streaming 是数据边界选择

批处理面对有界数据集:某天文件、某次快照、一个日期区间。流处理面对持续到来的无界事件,系统用窗口和 watermark 形成暂时可计算的边界。

维度BatchStreaming
结果延迟分钟到天毫秒到分钟
输入边界明确结束持续到达
失败恢复重跑批次checkpoint + replay
迟到数据下一批回补watermark/allowed lateness
运营复杂度通常较低状态、顺序与背压更复杂

同一系统可以“流式摄取、批量分析”,也可以用微批实现近实时。业务需要五分钟内告警,不意味着全部历史分析都必须改成流处理。

2. 摄取先建立批次身份与清单

读取每日 CSV 时,至少记录:

text
batch_id
source_uri
source_version / etag
expected_size
checksum
observed_rows
schema_version
started_at / completed_at
status
python
from 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. 分层存储隔离责任

一种常见但非唯一的分层:

text
raw        保留接收到的原始表示与 provenance
validated  通过 schema/基础质量检查,坏记录隔离
curated    语义统一、去重、业务规则明确
serving    针对报表、API 或模型组织

每层应有明确输入、输出和重建方式。不要把“bronze/silver/gold”当作固定真理;名称不重要,可重放边界与质量承诺才重要。

存储选择由访问模式、更新语义、规模、延迟和成本共同决定。列式文件适合扫描分析,但小文件过多会拖累元数据和调度;行式数据库适合点查与事务,却不意味着不能聚合。不要仅凭“按天聚合”就断言必须换某种存储。

5. 处理步骤必须幂等

任务在写到一半时失败,重跑不能把结果重复累加。常见策略:

  • 按分区写临时位置,成功后原子发布;
  • 以业务键执行 upsert;
  • 输出路径包含确定性 run/version;
  • checkpoint 与外部副作用协调;
  • 记录输入版本到输出版本的 lineage。

“先删目标分区再重写”虽然简单,却会产生暂时空窗,并在删除后失败时丢失旧结果。更安全的是写新版本、校验、再切换元数据指针。

6. Backfill 与日常运行使用同一逻辑

修正历史 bug 时需要回填。若回填另写临时脚本,规则很快漂移。管道应参数化:

text
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. 删除必须覆盖派生与备份

删除地图包括:

text
原始对象 → 清洗表 → 聚合表 → 特征 → 搜索索引 → 缓存 → 备份

删除文件名不等于物理介质立即不可恢复。现代 SSD、对象存储和托管数据库的擦除由底层实现决定,反复覆写未必可靠或可用。常见方法是:

  • 逻辑删除后由存储生命周期清除版本;
  • crypto erasure:销毁独立加密密钥;
  • 使用供应商保证的介质清退与擦除流程;
  • 对备份设置不可变但有限的过期窗口;
  • 保存不可含敏感内容的删除审计证据。

法律上需要删除什么、何时完成取决于适用规则和角色,应由合规与安全团队确认,不能由教程给出统一期限。

常见误区

  • 任务显示成功就代表数据完整:成功状态不覆盖语义核对。
  • 行数相等就没有丢失:重复可以抵消缺失。
  • 原始层永远保留最安全:会扩大敏感数据与合规风险。
  • rm 或反复覆写适用于所有存储:介质、版本和托管层行为不同。

练习

  1. 为每日快照设计 batch manifest 和重跑幂等策略。
  2. 构造“行数相等但一丢一重”的核对反例。
  3. 为一个历史回填写发布、验证与回滚步骤。
  4. 画出用户删除请求需要覆盖的所有派生数据节点。

小结

可运营的数据管道要能说明输入版本、处理状态、输出版本和失败恢复方式。批流选择决定边界,幂等与重放处理故障,质量监控判断结果是否可信,保留与删除则控制数据在何时退出系统。

下一章进入清洗工作台,但不会从“空值填中位数”开始,而是先识别缺失机制、类型歧义和每条修复规则正在做出的假设。

Built with VitePress | Software Systems Atlas