小项目:构建一个数据准备 CLI
把前面的迁移知识组合成一个可测试、可扩展的数据准备命令行工具。
学习目标
这是 Python/AI 专项的收束项目。完成后你能够:
- 用 argparse、dataclass 和纯函数组织一个可测试的数据准备 CLI。
- 对 JSONL 事件做 schema 校验、缺失处理、重复检测、特征化和统计。
- 输出稳定的特征文件与 manifest,记录数据版本、配置、哈希、保留和拒绝原因。
- 用 dry-run、错误退出码、原子写入和固定排序,让管线可重跑、可审核、可接入 CI。
从 JS/TS 迁移的心智模型
JavaScript/TypeScript 脚本可能是 rawEvents.filter(isValid).map(toFeature),最后写一个 JSON 文件。AI 数据准备不能把“过滤成功”当作交付:要知道输入的每一行是否解析,为什么被拒绝,特征的单位和标签是什么,输出能否被训练再次读取。Python CLI 把 IO、配置、业务规则和退出码分开,后面换成本地文件、对象存储或服务输入时,纯函数仍然可测试。
const result = rawEvents
.filter(isValid)
.map(toFeature);
await fs.promises.writeFile(output, JSON.stringify(result)); config = parse_args()
stats = prepare(config)
write_manifest(stats, config)
return 0 if stats["errors"] == 0 else 2 项目边界和输入输出契约
输入采用 UTF-8 JSONL,每行是一个事件对象,至少包含 event_id、device_id、timestamp、value、label。契约还要规定 timestamp 的时区、value 的单位和范围、label 的允许集合、文本最大长度、重复 event_id 的策略。不要用一个模糊的 dict 让下游猜字段;读取、解析、校验和特征化都应返回可追踪的结果。
输出 JSONL 每行包含 event_id、device_id、feature_value、label、schema_version 和 data_version。manifest 记录 input_sha256、output_sha256、config_digest、code_version、input_rows、kept_rows、dropped_rows、dropped_by_reason、feature_names 和运行时间。原始文本若敏感,输出只保留必要特征或哈希,不能把完整请求写进日志。
示例一:纯函数校验一行事件
from datetime import datetime
import math
ALLOWED_LABELS = {"wake", "noise"}
def parse_event(row, line_number, data_version):
if not isinstance(row, dict):
return None, {"line": line_number, "reason": "not_object"}
event_id = row.get("event_id")
device_id = row.get("device_id")
value = row.get("value")
label = row.get("label")
try:
timestamp = datetime.fromisoformat(str(row["timestamp"]).replace("Z", "+00:00"))
except (KeyError, TypeError, ValueError):
return None, {"line": line_number, "reason": "timestamp_invalid"}
if not event_id or not device_id:
return None, {"line": line_number, "reason": "id_missing"}
if not isinstance(value, (int, float)) or not math.isfinite(value):
return None, {"line": line_number, "reason": "value_invalid"}
if label not in ALLOWED_LABELS:
return None, {"line": line_number, "reason": "label_unknown"}
return {
"event_id": str(event_id),
"device_id": str(device_id),
"timestamp": timestamp.isoformat(),
"feature_value": float(value),
"label": label,
"data_version": data_version,
}, None
运行时对每一行返回 accepted 或带 line、reason 的 rejected,错误不会污染下一行。datetime 解析后统一时区,value 变成有限 float,label 通过 allowlist。这里没有用未来事件或 label 派生特征;若要计算设备均值,必须在切分和时间窗口规则中明确 fit 范围,避免 data leakage。
配置、命令行和确定性
CLI 的入口只负责读取参数、解析配置和调用 prepare。用 dataclass 表示 resolved config,固定默认值和覆盖顺序;输入输出路径相同应立即拒绝。dry-run 只读取、校验和统计,不写特征文件,但可以写一个明确标记为 dry-run 的报告。limit 是调试限制,不应悄悄成为生产默认值。
稳定输出需要固定读取顺序、错误原因排序、字段顺序和 JSON 序列化格式。不要把当前时间或随机 UUID 放进影响内容哈希的字段;运行时间可以单独写入审计记录。配置包含 seed 时也要让后续抽样显式使用它,并将 seed 写进 manifest。
示例二:用 dataclass 和 argparse 解析配置
import argparse
from dataclasses import dataclass
from pathlib import Path
@dataclass(frozen=True)
class Config:
input: Path
output: Path
manifest: Path
allowed_labels: tuple[str, ...]
data_version: str
schema_version: str
limit: int | None = None
dry_run: bool = False
def validate(self):
if self.input.resolve() == self.output.resolve():
raise ValueError("input and output must be different")
if not self.allowed_labels:
raise ValueError("at least one label is required")
if self.limit is not None and self.limit < 1:
raise ValueError("limit must be positive")
def parse_args(argv=None):
parser = argparse.ArgumentParser(prog="prepare-data")
parser.add_argument("--input", required=True, type=Path)
parser.add_argument("--output", required=True, type=Path)
parser.add_argument("--manifest", required=True, type=Path)
parser.add_argument("--label", action="append", required=True)
parser.add_argument("--data-version", required=True)
parser.add_argument("--schema-version", default="features-v1")
parser.add_argument("--limit", type=int)
parser.add_argument("--dry-run", action="store_true")
args = parser.parse_args(argv)
config = Config(
input=args.input, output=args.output, manifest=args.manifest,
allowed_labels=tuple(sorted(set(args.label))),
data_version=args.data_version,
schema_version=args.schema_version,
limit=args.limit, dry_run=args.dry_run,
)
config.validate()
return config
运行 prepare-data —help 应显示参数;缺少 —input、负 limit 或 input==output 时应非零退出,并给出可行动错误。配置解析成功后只传递 Config,不在各函数中读取环境变量。这样测试可以直接调用 parse_args 和 prepare,而不必启动子进程。
管线、统计和原子输出
prepare 的流程可以拆为 read_jsonl、parse_event、deduplicate、build_feature、write_jsonl、write_manifest。每层只有一个边界:读文件负责 UTF-8 和坏 JSON,校验负责字段,去重负责 event_id,特征函数负责数值表示,写出负责临时文件和原子替换。错误原因使用固定集合,例如 json_invalid、id_missing、timestamp_invalid、value_invalid、label_unknown、duplicate_id。
重复 ID 不能静默取最后一条,除非契约明确按时间排序后 keep-last。输出前按 timestamp、device_id、event_id 稳定排序;顺序稳定会让 output hash、训练样本和回归差异可比较。若输入文件很大,逐行处理并只保留统计,避免把全量 JSONL 读进内存;如果必须排序,明确外部排序或资源预算。
示例三:逐行处理并生成 manifest 统计
import hashlib
import json
from collections import Counter
def iter_events(path, config):
with path.open("r", encoding="utf-8") as handle:
for line_number, line in enumerate(handle, start=1):
if config.limit is not None and line_number > config.limit:
break
try:
row = json.loads(line)
except json.JSONDecodeError:
yield None, {"line": line_number, "reason": "json_invalid"}
continue
yield parse_event(row, line_number, config.data_version)
def digest_file(path):
digest = hashlib.sha256()
with path.open("rb") as handle:
for block in iter(lambda: handle.read(1024 * 1024), b""):
digest.update(block)
return digest.hexdigest()
def prepare(config):
kept, rejected = [], Counter()
seen = set()
for event, error in iter_events(config.input, config):
if error:
rejected[error["reason"]] += 1
continue
if event["event_id"] in seen:
rejected["duplicate_id"] += 1
continue
seen.add(event["event_id"])
kept.append(event)
kept.sort(key=lambda item: (item["timestamp"], item["device_id"], item["event_id"]))
return {
"input_rows": sum(rejected.values()) + len(kept),
"kept_rows": len(kept),
"dropped_rows": sum(rejected.values()),
"dropped_by_reason": dict(sorted(rejected.items())),
"input_sha256": digest_file(config.input),
"records": kept,
}
这个示例为教学清晰度保留 kept 列表;真实大文件要将排序策略与内存预算写入设计。manifest 的 input_rows、kept_rows、dropped_rows 必须满足守恒关系;rejected reason 排序后输出,避免字典顺序造成无意义 hash 差异。任何目标派生字段都应在 split 后按训练范围计算,防止数据泄漏。
运行、输出与验证
用五行 fixture 验收:一行合法、一行坏 JSON、一行缺 ID、一行重复 ID、一行未知 label。执行:
python -m prepare_data --input events.jsonl --output features.jsonl --manifest run.json --label wake --label noise --data-version events-v3
python -m prepare_data --input events.jsonl --output features.jsonl --manifest dry-run.json --label wake --label noise --data-version events-v3 --dry-run
预期 stdout 或日志包含 input_rows=5、kept_rows=1、dropped_rows=4,并列出每类 reason;dry-run 不应创建 features.jsonl。运行两次时,记录的 output 内容和核心 manifest 哈希应一致,运行时间字段可以不同。用 JSON 解析器读取每一行,验证字段顺序、UTF-8、label allowlist、feature shape 和有限值。
CLI 的退出码也是输出契约:正常运行返回 0,存在契约级错误或写入失败返回非零;被拒绝的坏行是否让整个命令失败,要由配置决定,但必须在 manifest 中计数。输入不存在、输出目录无权限、磁盘不足和 input==output 属于命令失败,不要把它们混成 dropped row。
常见错误、排错与调试
- kept 数量与输入行数不守恒:检查坏 JSON 是否计数、limit 是否生效、重复行是否被重复统计。
- 两次 output hash 不同:比较输入 hash、稳定排序、字段格式、浮点序列化和配置 digest,先排除时间字段。
- 模型指标突然变化:比较 data_version、schema_version、feature_names、丢弃原因和 train/valid/test split;不要立刻调网络。
- 训练数据泄漏:检查是否用全量未来统计量、label 派生字段或同一 device 的跨 split 聚合;将 fit 范围写入 manifest。
- CLI 在 Windows 编码失败:明确使用 UTF-8,路径用 pathlib,错误文件保留 line number;不要依赖当前工作目录。
- 输出损坏:写到 output.tmp,flush 后原子 replace;重复运行时不要在原文件上半写。
- dry-run 仍改文件:测试所有写函数入口,dry_run 只允许报告输出;用临时目录比较目录快照。
- 处理大文件 OOM 或超时:逐行读取、限制日志、测量读写和排序时间;必要时做分块或外部排序并记录资源。
- 日志暴露隐私:只记录 ID 哈希、shape、数量和 reason,不打印 text、token、API key 或完整原始行。
练习与任务
实现完整 prepare-data:读取 JSONL,按 schema 校验并去重,生成 AI 特征 JSONL 和 manifest;支持 label、data-version、schema-version、limit、dry-run;输入输出不可相同,输出原子写入,错误有行号和固定 reason。写 pytest 覆盖正常、坏 JSON、缺字段、NaN、重复 ID、未知标签、空输入、重复运行和非零退出码。
可审核数据准备 CLI
完成 parse_args、parse_event、prepare、write_outputs 和 main(argv),返回 accepted/rejected 统计;manifest 要包含输入哈希、配置摘要、schema/data version、丢弃原因和稳定输出哈希。
给我一点提示
先定义 Config 和 reason 集合;纯函数先测,再接 pathlib、argparse 和原子文件;dry-run 不写特征文件。
查看参考答案
config = parse_args(argv)
stats = prepare(config)
if not config.dry_run:
write_outputs(stats, config)
return 0 if stats["dropped_rows"] == 0 else 2 完整答案
def write_outputs(stats, config):
config.output.parent.mkdir(parents=True, exist_ok=True)
config.manifest.parent.mkdir(parents=True, exist_ok=True)
temporary = config.output.with_suffix(config.output.suffix + ".tmp")
with temporary.open("w", encoding="utf-8", newline="\n") as handle:
for record in stats["records"]:
handle.write(json.dumps(record, ensure_ascii=False, sort_keys=True) + "\n")
temporary.replace(config.output)
manifest = {
"schema_version": config.schema_version,
"data_version": config.data_version,
"input_sha256": stats["input_sha256"],
"input_rows": stats["input_rows"],
"kept_rows": stats["kept_rows"],
"dropped_rows": stats["dropped_rows"],
"dropped_by_reason": stats["dropped_by_reason"],
"config": {
"labels": config.allowed_labels,
"limit": config.limit,
"dry_run": config.dry_run,
},
"output_sha256": digest_file(config.output),
}
config.manifest.write_text(
json.dumps(manifest, ensure_ascii=False, indent=2, sort_keys=True),
encoding="utf-8",
)
def main(argv=None):
try:
config = parse_args(argv)
stats = prepare(config)
if not config.dry_run:
write_outputs(stats, config)
print(json.dumps({
"input_rows": stats["input_rows"],
"kept_rows": stats["kept_rows"],
"dropped_rows": stats["dropped_rows"],
"dropped_by_reason": stats["dropped_by_reason"],
}, ensure_ascii=False, sort_keys=True))
return 0 if stats["dropped_rows"] == 0 else 2
except (OSError, ValueError) as error:
print(f"prepare-data failed: {error}", file=sys.stderr)
return 1
main 需要导入 sys,并在模块底部使用 SystemExit(main())。这里选择“有坏行则返回 2”,也可以改成允许部分成功但必须保持协议一致。用相同 fixture 重跑,比较 features.jsonl 的输出哈希;用 dry-run 比较目录快照,确认没有产生数据文件;再用输入输出同路径、目录不存在和磁盘错误模拟失败路径。
本节结论
项目完成的证据是可重跑的输入、稳定的输出、可解释的拒绝统计、明确的退出码和能阻断泄漏的契约。CLI 只是入口,真正可复用的是纯校验、特征化和 manifest 边界。
与同一 AI 项目主线的连接
这条 CLI 复用了 pandas 的字段清洗思想、NumPy 的 dtype/shape 契约、数据集课程的 split 和 leakage 审计、PyTorch 的 Dataset/评估记录以及 MLOps 的 data_version 和 manifest。它输出的特征文件可以进入 baseline 或 DataLoader,manifest 可以随 checkpoint 和 deployment 一起保存。发生指标变化时,先比较输入哈希、schema、丢弃原因和特征顺序,再判断模型是否退化;上线服务也能用同一 schema 做请求校验。
小结
一个 AI 数据准备项目不需要很多框架,但必须把输入契约、纯函数、配置、统计、版本、稳定排序、原子输出、退出码和测试组合起来。JSONL 中每一行都要可追踪,坏数据要有原因,dry-run 要安全,重复运行要可比较,未来统计要遵守 split 边界。完成这个 CLI,你就把 JS/TS 的脚本式数据转换升级成了可复现、可审核的 Python AI 工程闭环。
延伸阅读
先完成本节练习,再用这些资料查阅完整 API 和真实项目组织方式。
本节是阶段检查点。完成练习后,再进入下一阶段。