3.4 仓储管理:数据存储与组织 本节摘要:采集数据存储要回答三个问题——原始层要不要留、按什么分区、版本怎么管。本节用"三层数仓加版本血缘"的组织法落地:原始层只进不改、加工层按日期与来源分区、成品层带版本号与血缘记录,并给出 JSONL 加 Parquet 的实现与检索层的取舍。 仓储的三个问题 清洗车间出来的合格数据往哪放,多数团队的第一反应是"插数据库"。真正的问题藏在三个追问里。原始层要不要留? 要。清洗规则永远会改——3.3 节那条语义规则上周刚误杀过一批数据,没有 raw 层的原始页面存档,重跑规则等于重新爬一遍,纪律与预算都遭殃。按什么分区? 按访问模式定:增量任务按日期查、多源任务按来源过滤、质量问题按批次追溯,分区键就应该是这三个查询的 WHERE 条件。版本怎么管?
本节摘要:采集数据存储要回答三个问题——原始层要不要留、按什么分区、版本怎么管。本节用"三层数仓加版本血缘"的组织法落地:原始层只进不改、加工层按日期与来源分区、成品层带版本号与血缘记录,并给出 JSONL 加 Parquet 的实现与检索层的取舍。
清洗车间出来的合格数据往哪放,多数团队的第一反应是"插数据库"。真正的问题藏在三个追问里。原始层要不要留? 要。清洗规则永远会改——3.3 节那条语义规则上周刚误杀过一批数据,没有 raw 层的原始页面存档,重跑规则等于重新爬一遍,纪律与预算都遭殃。按什么分区? 按访问模式定:增量任务按日期查、多源任务按来源过滤、质量问题按批次追溯,分区键就应该是这三个查询的 WHERE 条件。版本怎么管? 清洗规则改版、数据源增减都会产出"新版数据集",模型侧永远需要回答"这个结果是哪版数据训的"——没有版本与血缘,这个问题无解。
三个问题合成一个组织法:三层数仓加版本血缘。raw 层存原始抓取产物(HTML 快照或原始 JSON),只进不改;staging 层存清洗中间态,按日期加来源分区,可重算;curated 层存对齐备料单的成品,带版本号,供训练与检索直接消费。

不引入重型数据平台,用文件系统加约定就能把组织法落地。分区目录结构是核心,血缘用一份 JSON 描述:
import json, hashlib from pathlib import Path from datetime import date class DatasetWarehouse: """最小三层仓储:分区写入加血缘登记""" def __init__(self, root: str): self.root = Path(root) def write_staging(self, rows: list[dict], source: str, day: date) -> str: """staging 层按 日期/来源 分区写 JSONL""" part = self.root / "staging" / f"dt={day:%Y%m%d}" / f"src={source}" part.mkdir(parents=True, exist_ok=True) path = part / "part-0000.jsonl" with path.open("a", encoding="utf-8") as f: # 追加写:增量任务不覆盖历史 for r in rows: f.write(json.dumps(r, ensure_ascii=False) + "\n") return str(path.relative_to(self.root)) def publish_curated(self, rows: list[dict], version: str, lineage: dict) -> str: """curated 层发布带血缘的版本目录""" vdir = self.root / "curated" / f"v{version}" vdir.mkdir(parents=True, exist_ok=True) (vdir / "data.jsonl").write_text( "\n".join(json.dumps(r, ensure_ascii=False) for r in rows), encoding="utf-8") (vdir / "lineage.json").write_text( json.dumps(lineage, ensure_ascii=False, indent=2), encoding="utf-8") return str(vdir.relative_to(self.root)) wh = DatasetWarehouse(root="warehouse") p = wh.write_staging([{"q": "时效几年", "a": "一年"}], source="legal_example", day=date(2026, 8, 26)) print(p) # 输出:staging/dt=20260826/src=legal_example/part-0000.jsonl lineage = { "version": "1.3.0", "from_partitions": ["staging/dt=20260801至20260826/src=legal_example"], "clean_rules": "clean@2.1.0(字符层NFKC 结构层日期ISO 语义层len大于5)", "stats": {"total": 51230, "completeness": 0.982, "duplicate_ratio": 0.031}, "published_at": "2026-08-26", "owner": "data-team", } v = wh.publish_curated([{"q": "时效几年", "a": "一年"}], version="1.3.0", lineage=lineage) print(v) # 输出:curated/v1.3.0
日期与来源组成的双分区键,价值在两个场景兑现:增量任务只读今天的分区,扫描量与全量无关;某来源暴出假数据(2.4 节的蜜罐场景),按来源键隔离该来源、其余照常服役。血缘文件里的 clean_rules 字段带版本号——规则改版时递增版本号,旧版本数据集的血缘能精确指向旧规则,复现链路闭合。
文件仓储之外,三类常见消费者各有对口存储。训练管线批量顺序读,Parquet 列存按列裁剪、体积小加载快,是 curated 层的首选格式;检索服务要语义近邻查询,向量库(Milvus、Chroma 一族)对口,embedding 生成管线接在 curated 层下游;业务后台要交互式查询与事务,关系库(PostgreSQL 一族)对口,通常只入统计与样本,不入全量。让文件层做数据湖底座、三类库做消费视图,比"一个库装一切"清爽得多。
常见坑:raw 层顺手做了"一点小清洗"(去个空白、转个码)。第一次重跑规则时就会发现:原始数据已被污染,"重算"出来的结果和线上对不上。raw 层的纪律是只进不改,连编码都保持原样,解码留到读取时。
关键直觉:仓储的价值在灾难时兑现。清洗规则误杀、数据源暴雷、模型效果回退——三种事故的恢复动作都是"从 raw 重算到指定版本",这套动作能多快执行完,就是仓储设计的成绩单。