【技术专栏】 AI开发
【内容摘要】 综合运用批量处理、结构化输出、结果校验与幂等写入,把一批文本安全地整理成可回滚的标签数据。
前面的系列已经覆盖了模型调用、提示词、结构化输出、错误处理、评估和安全边界。本篇在路线完成后做一个新的综合小项目:读取一组带有唯一编号的文本,请模型为每条文本选择已有标签,程序校验结果并写入新文件。它不修改原始数据,也不把模型输出直接覆盖到生产系统,重点是练习批处理中的可追踪、可重试和可回滚。
先定义数据和边界
输入使用 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 负责数据完整性和副作用边界,离线测试负责保护确定性代码。真正扩展到生产时,还需要加入有限重试、速率控制、成本统计和人工抽样,但每增加一个能力,都应先定义失败状态和恢复方式。
技术分享自:时光笔记 (wxy.email) | 华为云开发者社区
【声明】本内容来自华为云开发者社区博主,不代表华为云及华为云开发者社区的观点和立场。转载时必须标注文章的来源(华为云社区)、文章链接、文章作者等基本信息,否则作者和本社区有权追究责任。如果您发现本社区中有涉嫌抄袭的内容,欢迎发送邮件进行举报,并提供相关证据,一经查实,本社区将立刻删除涉嫌侵权内容,举报邮箱:
cloudbbs@huaweicloud.com
评论(0)