Python / AI · 服务化与结课项目 · LESSON 30

小项目:构建一个数据准备 CLI

把前面的迁移知识组合成一个可测试、可扩展的数据准备命令行工具。

18 分钟project · cli · ai · testing

学习目标

这是 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、配置、业务规则和退出码分开,后面换成本地文件、对象存储或服务输入时,纯函数仍然可测试。

TRANSLATION LENS 同一个意图,两种工程表达 窄屏可左右滑动查看完整代码
JS / TS
const result = rawEvents
.filter(isValid)
.map(toFeature);
await fs.promises.writeFile(output, JSON.stringify(result));
Python CLI
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、未知标签、空输入、重复运行和非零退出码。

01
TRY IT YOURSELF

可审核数据准备 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 工程闭环。

FURTHER READING

延伸阅读

先完成本节练习,再用这些资料查阅完整 API 和真实项目组织方式。

当前学习阶段服务化与结课项目
0/6

本节是阶段检查点。完成练习后,再进入下一阶段。