前面的系列已经覆盖了模型调用、提示词、结构化输出、错误处理、评估和安全边界。本篇在路线完成后做一个新的综合小项目:读取一组带有唯一编号的文本,请模型为每条文本选择已有标签,程序校验结果并写入新文件。它不修改原始数据,也不把模型输出直接覆盖到生产系统,重点是练习批处理中的可追踪、可重试和可回滚。

先定义数据和边界

输入使用 JSON Lines(JSONL),每行一个对象,只需要 id 和 text。输出仍然保留 id、原文和 tags,其中标签只能从程序给出的白名单中选择。模型可以判断文本属于哪些标签,但不能新造标签、改写原文或改变编号。

程序分成五步:读取并检查输入、逐条调用模型、解析和校验结果、写入临时文件、原子替换输出文件。原始文件始终只读;处理中断时,临时文件不会冒充成功结果。输出文件使用新的路径,因此删除输出文件即可回滚到原始数据。

准备环境和输入

创建虚拟环境并安装官方 SDK 和环境变量加载库:

1
2
3
python -m venv .venv
source .venv/bin/activate
python -m pip install openai python-dotenv

创建 .env,密钥只能通过环境变量提供:

1
2
3
OPENAI_API_KEY=替换为你的真实密钥
MODEL_NAME=替换为你可用的模型名称
OPENAI_BASE_URL=

不要把 .env 提交到 Git。准备 articles.jsonl,例如每行一个对象:

1
2
{"id":"a-001","text":"介绍如何用 Python 读取环境变量并隔离 API 密钥。"}
{"id":"a-002","text":"记录一次向量检索的召回率测试方法。"}

标签白名单由程序维护,例如 python、security、rag、testing。白名单是业务约束,不应只写在提示词里。

编写批处理程序

新建 batch_tagger.py:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
import json
import os
import sys
import tempfile
from pathlib import Path

from dotenv import load_dotenv
from openai import OpenAI

load_dotenv()
MODEL = os.environ["MODEL_NAME"]
client = OpenAI(
api_key=os.environ.get("OPENAI_API_KEY"),
base_url=os.environ.get("OPENAI_BASE_URL") or None,
)
ALLOWED_TAGS = {"python", "security", "rag", "testing"}
RULES = """
根据输入文本选择相关标签。只返回 JSON 对象,字段只能是 tags。
tags 必须是数组,元素只能来自 python、security、rag、testing,最多选择三个;
没有匹配项时返回空数组。不要改写文本,不要解释原因。
""".strip()


def validate_item(raw: str, item_id: str) -> dict:
try:
value = json.loads(raw)
except json.JSONDecodeError as exc:
raise ValueError(f"{item_id}: 返回不是合法 JSON") from exc
if set(value) != {"tags"} or not isinstance(value["tags"], list):
raise ValueError(f"{item_id}: 输出字段无效")
tags = value["tags"]
if len(tags) > 3 or len(set(tags)) != len(tags):
raise ValueError(f"{item_id}: 标签数量或重复项无效")
if any(not isinstance(tag, str) or tag not in ALLOWED_TAGS for tag in tags):
raise ValueError(f"{item_id}: 出现白名单之外的标签")
return {"id": item_id, "tags": tags}


def classify(text: str, item_id: str) -> dict:
if not text.strip() or len(text) > 6000:
raise ValueError(f"{item_id}: 文本为空或超过长度限制")
response = client.responses.create(
model=MODEL,
instructions=RULES,
input=text,
)
return validate_item(response.output_text, item_id)


def read_input(path: Path) -> list[dict]:
records = []
seen = set()
for line_no, line in enumerate(path.read_text(encoding="utf-8").splitlines(), 1):
if not line.strip():
continue
try:
record = json.loads(line)
except json.JSONDecodeError as exc:
raise ValueError(f"第 {line_no} 行不是合法 JSON") from exc
if set(record) != {"id", "text"} or not isinstance(record["id"], str):
raise ValueError(f"第 {line_no} 行字段无效")
if record["id"] in seen:
raise ValueError(f"第 {line_no} 行存在重复 id")
seen.add(record["id"])
records.append(record)
return records


def process(source: Path, destination: Path) -> None:
records = read_input(source)
results = []
for record in records:
tagged = classify(record["text"], record["id"])
results.append({**record, **tagged})

destination.parent.mkdir(parents=True, exist_ok=True)
with tempfile.NamedTemporaryFile(
"w", encoding="utf-8", dir=destination.parent, delete=False
) as handle:
temporary = Path(handle.name)
for result in results:
handle.write(json.dumps(result, ensure_ascii=False) + "\n")
temporary.replace(destination)


def main() -> None:
if len(sys.argv) != 3:
raise SystemExit("用法:python batch_tagger.py articles.jsonl tagged.jsonl")
process(Path(sys.argv[1]), Path(sys.argv[2]))
print(f"已写入 {sys.argv[2]}")


if __name__ == "__main__":
main()

运行 python batch_tagger.py articles.jsonl tagged.jsonl。responses.create 负责单条请求,output_text 取出文本结果;真正写文件前,程序已经检查了 JSON、字段集合、标签白名单、数量和重复项。代码中的 temporary.replace(destination) 会在同一文件系统内替换目标文件,处理失败时不会留下一个看似完整的半成品。

为什么批处理不能直接覆盖

如果循环中每处理一条就打开目标文件写入,程序在第 5 条失败时,前 4 条可能已经覆盖旧结果。下一次运行还无法判断哪些条目成功过。这里选择“全部成功后再替换”,虽然失败时会丢弃本轮已得到的模型结果,却换来了清晰的成功和失败状态。

如果数据量很大,可以增加检查点文件,但检查点必须记录输入版本、id、提示词版本和模型名称。重试时只处理没有成功的编号,最后再按照原始顺序合并。不要仅凭行号续跑,因为输入一旦插入或删除记录,行号就会变化。

用离线测试保护确定性逻辑

模型调用不稳定,也会产生费用;先测试校验函数。创建 test_batch_tagger.py:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
import json
import unittest

from batch_tagger import validate_item


class TaggerTest(unittest.TestCase):
def test_valid_tags(self):
raw = json.dumps({"tags": ["python", "testing"]})
self.assertEqual(validate_item(raw, "a-001")["tags"], ["python", "testing"])

def test_unknown_tag_is_rejected(self):
raw = json.dumps({"tags": ["python", "news"]})
with self.assertRaises(ValueError):
validate_item(raw, "a-002")

def test_extra_field_is_rejected(self):
raw = json.dumps({"tags": [], "reason": "unused"})
with self.assertRaises(ValueError):
validate_item(raw, "a-003")


if __name__ == "__main__":
unittest.main()

运行 python -m unittest -v test_batch_tagger.py。这些测试不访问网络,验证的是“坏结果不会进入写入阶段”。还应为重复标签、超过三个标签、重复 id、空文本和非法输入行增加测试。至于标签是否符合文章真实含义,属于模型质量和业务评估问题,需要准备带人工标注的样本集,不能靠类型校验解决。

常见问题

模型每次选择的标签不一样怎么办? 先固定白名单和输出约束,再用标注样本比较准确率、漏标率和多标率。批处理应保存运行时间、模型名和提示词版本,方便定位变化来源。

单条失败时要不要继续? 这取决于业务。如果结果必须完整,当前程序会直接失败并不替换文件;如果允许部分成功,可以记录失败的 id,但输出中必须明确标记未处理项,不能把缺失当成空标签。

为什么不自动创建新标签? 新标签会改变下游过滤、统计和权限边界。可以另写一个候选报告,由人工审核后再更新白名单,而不是在批处理中悄悄扩展集合。

如何回滚? 原始 JSONL 从未被修改,删除或改名 tagged.jsonl 即可回到未整理状态。若以后接入数据库,应先保存批次号和旧值,再设计事务或反向迁移,而不是把文件替换思路直接套到所有系统。

小结

这个综合项目把模型调用放进了一个可回滚的批处理流程:输入有唯一编号,输出有明确白名单,模型结果经过严格校验,文件通过临时路径完成原子替换。模型负责提出标签,Python 负责数据完整性和副作用边界,离线测试负责保护确定性代码。真正扩展到生产时,还需要加入有限重试、速率控制、成本统计和人工抽样,但每增加一个能力,都应先定义失败状态和恢复方式。