前面的系列已经覆盖了模型调用、提示词、结构化输出、错误处理、评估和安全边界。本篇在路线完成后做一个新的综合小项目:读取一组带有唯一编号的文本,请模型为每条文本选择已有标签,程序校验结果并写入新文件。它不修改原始数据,也不把模型输出直接覆盖到生产系统,重点是练习批处理中的可追踪、可重试和可回滚。
先定义数据和边界 输入使用 JSON Lines(JSONL),每行一个对象,只需要 id 和 text。输出仍然保留 id、原文和 tags,其中标签只能从程序给出的白名单中选择。模型可以判断文本属于哪些标签,但不能新造标签、改写原文或改变编号。
程序分成五步:读取并检查输入、逐条调用模型、解析和校验结果、写入临时文件、原子替换输出文件。原始文件始终只读;处理中断时,临时文件不会冒充成功结果。输出文件使用新的路径,因此删除输出文件即可回滚到原始数据。
准备环境和输入 创建虚拟环境并安装官方 SDK 和环境变量加载库:
1 2 3 python -m venv .venv source .venv/bin/activatepython -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 jsonimport osimport sysimport tempfilefrom pathlib import Pathfrom dotenv import load_dotenvfrom openai import OpenAIload_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 jsonimport unittestfrom batch_tagger import validate_itemclass 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 负责数据完整性和副作用边界,离线测试负责保护确定性代码。真正扩展到生产时,还需要加入有限重试、速率控制、成本统计和人工抽样,但每增加一个能力,都应先定义失败状态和恢复方式。