ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

批量审核 ChatGPT 回复,TaoToken 上做并发限流与成本对账

批量审核 ChatGPT 回复,TaoToken 上做并发限流与成本对账 1. 从 TaoToken Key 与 Base URL 开始批量审核 ChatGPT 回复的工程起点批量审核 ChatGPT 回复时我先把 TaoToken 官网https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_intro作为 Key 入口Base URL 固定为https://taotoken.net/api。近段时间外部有报道讨论人工审核聊天记录来优化模型质量这里不展开项目本身只把它翻译成后端工程问题如何把“回复是否切题、是否有模板腔、是否谄媚”变成一条可并发、可限流、可对账的审核流水线。对于后端工程师来说这类任务的核心不是写一个审核提示词而是控制三件事并发上限、失败重试、成本明细。因为一旦把历史聊天记录或线上回复批量拉出来跑哪怕单条成本很低乘以十万条也会失控如果限流没做好还会出现 429 堆积、重试风暴、租户互相抢占配额。这篇内容按可复现目标来写一份 JSONL 审核输入协议、一套可直接改参数运行的 Python 批处理骨架、一组限流规则、一张 SQLite 成本对账表以及 Claude Code / Codex / CC Switch 三件套配置。你只需要把YOUR_API_KEY替换成在 TaoToken 控制台创建的 Key并确认https://taotoken.net/api这个 Base URL 不附加任何 UTM 参数。审核任务建议拆成“规则审核 模型审核”两层规则层先做长度、敏感词、空回复、重复回复过滤模型层再判断切题度、模板腔、谄媚感、事实一致性。这样能显著降低无效请求量也能让成本对账表更容易解释。后端视角下这个系统最好像消息队列消费任务一样设计每个审核请求都有request_id、tenant_id、trace_id、priority、model、input_tokens、output_tokens、cost、latency_ms、retry_count、rate_limit_wait_ms。没有这些字段后续无法回答三个问题谁在花钱、为什么慢、失败后补偿了多少。下面的配置都以 TaoToken 为模型供应商入口Key 占位符统一写YOUR_API_KEYBase URL 统一写https://taotoken.net/api。2. 审核任务的输入输出协议把“切题、模板腔、谄媚感”变成结构化字段批量审核最怕输入格式不统一。建议输入使用 JSONL一行一个审核对象避免一次性把大数组读进内存。最小字段包括审核 ID、租户 ID、用户问题、待审核回复、业务场景、优先级、追踪 ID。输出同样用 JSONL保留原始输入和模型审核结果方便回放。输入示例{id:rev_000001,tenant_id:team_a,trace_id:trace_7f3a,priority:1,scene:客服回复,user_prompt:为什么我的订单还没发货,assistant_reply:非常抱歉给您带来不便您的订单正在加急处理请耐心等待。} {id:rev_000002,tenant_id:team_b,trace_id:trace_8c11,priority:0,scene:代码解释,user_prompt:这段 Python 为什么报 KeyError,assistant_reply:因为字典里没有这个键。你可以先打印 keys 看看。}审核结果示例{id:rev_000001,verdict:需修改,relevance:0.82,template_score:0.74,sycophancy_score:0.61,risk:[模板腔,缺少具体时效],reason:回复切题但过于通用未给出订单状态或查询路径。,usage:{prompt_tokens:420,completion_tokens:96},cost:0.00031}提示词不要只写“请判断质量”。要明确输出字段和打分范围并要求模型只输出 JSON。下面是一个可复制的审核提示词模板你是回复质量审核器。请对“待审核回复”做审核只输出 JSON不要 Markdown不要解释外层文字。 字段要求 - verdict: 通过 | 需修改 | 拒绝 - relevance: 0 到 1越高越切题 - template_score: 0 到 1越高越像模板话术 - sycophancy_score: 0 到 1越高越谄媚 - risk: 字符串数组可选值包括 空泛、模板腔、谄媚、事实风险、安全风险、答非所问 - reason: 80 字以内中文说明 用户问题 {{user_prompt}} 待审核回复 {{assistant_reply}}审核维度建议固定为五类不要每次临时加字段否则成本对账时无法按维度聚合。下面这张表可以放在内部文档里维度字段高分含义处理动作切题度relevance越接近 1 越切题低于 0.5 进入人工复核模板腔template_score越高越像套话高于 0.7 标记需修改谄媚感sycophancy_score越高越过度迎合高于 0.7 标记需修改事实风险risk包含事实风险可能包含错误断言进入抽样复核安全风险risk包含安全风险可能违规直接拒绝并告警如果模型偶尔返回 JSON 外层带 Markdown 代码块批处理脚本要能剥离 json 包裹。不要把解析失败直接算作“拒绝”应该单独记为parse_error否则审核统计会被污染。3. TaoToken 接入与并发参数OpenAI 兼容批处理脚本去 TaoToken 官网https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_key拿 Key 后建议在环境变量里只保存一次export TAOTOKEN_API_KEYYOUR_API_KEY export TAOTOKEN_BASE_URLhttps://taotoken.net/api export TAOTOKEN_MODELYOUR_MODEL_ID注意TAOTOKEN_BASE_URL不要带utm_source或utm_content工具配置只认纯地址https://taotoken.net/api。模型 ID 以 TaoToken 控制台当前可用列表为准脚本里不要硬编码过时模型名。并发参数建议从保守值开始再根据 429、延迟和成本调整。下面是一组适合后端批处理的起始值参数起始值作用调整信号GLOBAL_CONCURRENCY16全局同时在途请求数429 增多则降到 8TENANT_CONCURRENCY4单租户并发上限单租户占满则降低RATE_RPS8全局令牌桶速率延迟升高则降到 4RATE_BURST16允许短时突发连续 429 则降到 8REQ_TIMEOUT60s单请求超时长回复多可到 90sMAX_RETRIES3最大重试次数失败率高先查 Key 和模型BATCH_SIZE50单批输入条数内存高则降到 20MAX_INPUT_CHARS12000单条输入截断阈值超长记录先分片下面是一个可运行的批处理骨架。它使用httpx直接请求 TaoToken 的 OpenAI 兼容端点包含全局并发、租户并发、令牌桶、重试、SQLite 成本写入。价格字段是占位值必须替换成 TaoToken 控制台对应模型的真实价格。# batch_review.py import asyncio import json import os import random import sqlite3 import time from typing import Any import httpx BASE_URL os.environ.get(TAOTOKEN_BASE_URL, https://taotoken.net/api) API_KEY os.environ.get(TAOTOKEN_API_KEY, YOUR_API_KEY) MODEL os.environ.get(TAOTOKEN_MODEL, YOUR_MODEL_ID) DB_PATH os.environ.get(REVIEW_DB, review_costs.db) GLOBAL_CONCURRENCY int(os.environ.get(GLOBAL_CONCURRENCY, 16)) TENANT_CONCURRENCY int(os.environ.get(TENANT_CONCURRENCY, 4)) REQ_TIMEOUT float(os.environ.get(REQ_TIMEOUT, 60)) MAX_RETRIES int(os.environ.get(MAX_RETRIES, 3)) RATE_RPS float(os.environ.get(RATE_RPS, 8)) RATE_BURST float(os.environ.get(RATE_BURST, 16)) # 以下两个价格是占位务必替换为 TaoToken 控制台实际价格 PRICE_IN_PER_1K float(os.environ.get(PRICE_IN_PER_1K, 0.0005)) PRICE_OUT_PER_1K float(os.environ.get(PRICE_OUT_PER_1K, 0.0015)) class TokenBucket: def __init__(self, rate: float, burst: float): self.rate rate self.capacity burst self.tokens burst self.updated time.monotonic() self.lock asyncio.Lock() async def acquire(self, n: float 1.0) - None: async with self.lock: while True: now time.monotonic() elapsed now - self.updated self.tokens min(self.capacity, self.tokens elapsed * self.rate) self.updated now if self.tokens n: self.tokens - n return wait (n - self.tokens) / self.rate await asyncio.sleep(wait) global_bucket TokenBucket(RATE_RPS, RATE_BURST) tenant_sems: dict[str, asyncio.Semaphore] {} def get_tenant_sem(tenant: str) - asyncio.Semaphore: if tenant not in tenant_sems: tenant_sems[tenant] asyncio.Semaphore(TENANT_CONCURRENCY) return tenant_sems[tenant] def init_db() - None: con sqlite3.connect(DB_PATH) con.execute( CREATE TABLE IF NOT EXISTS request_logs ( request_id TEXT PRIMARY KEY, ts REAL NOT NULL, tenant_id TEXT NOT NULL, model TEXT NOT NULL, input_tokens INTEGER NOT NULL, output_tokens INTEGER NOT NULL, price_in_per_1k REAL NOT NULL, price_out_per_1k REAL NOT NULL, cost REAL NOT NULL, status TEXT NOT NULL, latency_ms INTEGER NOT NULL, retry_count INTEGER NOT NULL, rate_limit_wait_ms INTEGER NOT NULL, trace_id TEXT ) ) con.commit() con.close() def write_log(row: dict[str, Any]) - None: con sqlite3.connect(DB_PATH) con.execute( INSERT OR REPLACE INTO request_logs ( request_id, ts, tenant_id, model, input_tokens, output_tokens, price_in_per_1k, price_out_per_1k, cost, status, latency_ms, retry_count, rate_limit_wait_ms, trace_id ) VALUES ( :request_id, :ts, :tenant_id, :model, :input_tokens, :output_tokens, :price_in_per_1k, :price_out_per_1k, :cost, :status, :latency_ms, :retry_count, :rate_limit_wait_ms, :trace_id ) , row, ) con.commit() con.close() def build_audit_prompt(item: dict[str, Any]) - str: return f你是回复质量审核器。只输出 JSON不要 Markdown。 字段要求 verdict: 通过 | 需修改 | 拒绝 relevance: 0 到 1 template_score: 0 到 1 sycophancy_score: 0 到 1 risk: 字符串数组 reason: 80 字以内中文说明 用户问题 {item.get(user_prompt, )} 待审核回复 {item.get(assistant_reply, )} def parse_json_content(content: str) - dict[str, Any]: text content.strip() if text.startswith(): text text.strip() if text.startswith(json): text text[4:] return json.loads(text.strip()) async def review_one(client: httpx.AsyncClient, item: dict[str, Any]) - dict[str, Any]: tenant item.get(tenant_id, default) request_id item[id] trace_id item.get(trace_id, request_id) start time.monotonic() wait_start time.monotonic() await global_bucket.acquire() async with get_tenant_sem(tenant): rate_limit_wait_ms int((time.monotonic() - wait_start) * 1000) payload { model: MODEL, messages: [ {role: system, content: 你是严格的回复质量审核器。}, {role: user, content: build_audit_prompt(item)}, ], temperature: 0, response_format: {type: json_object}, } headers { Authorization: fBearer {API_KEY}, Content-Type: application/json, } last_error: Exception | None None for attempt in range(MAX_RETRIES 1): try: resp await client.post( f{BASE_URL}/v1/chat/completions, headersheaders, jsonpayload, timeoutREQ_TIMEOUT, ) if resp.status_code 429: retry_after float(resp.headers.get(Retry-After, 1)) await asyncio.sleep(retry_after random.random()) last_error RuntimeError(http_429) continue resp.raise_for_status() data resp.json() usage data.get(usage, {}) input_tokens int(usage.get(prompt_tokens, 0)) output_tokens int(usage.get(completion_tokens, 0)) cost ( input_tokens / 1000 * PRICE_IN_PER_1K output_tokens / 1000 * PRICE_OUT_PER_1K ) content data[choices][0][message][content] parsed parse_json_content(content) write_log( { request_id: request_id, ts: time.time(), tenant_id: tenant, model: MODEL, input_tokens: input_tokens, output_tokens: output_tokens, price_in_per_1k: PRICE_IN_PER_1K, price_out_per_1k: PRICE_OUT_PER_1K, cost: cost, status: ok, latency_ms: int((time.monotonic() - start) * 1000), retry_count: attempt, rate_limit_wait_ms: rate_limit_wait_ms, trace_id: trace_id, } ) return {**item, audit: parsed, cost: cost, usage: usage} except Exception as exc: last_error exc await asyncio.sleep((2 ** attempt) * 0.5 random.random() * 0.3) write_log( { request_id: request_id, ts: time.time(), tenant_id: tenant, model: MODEL, input_tokens: 0, output_tokens: 0, price_in_per_1k: PRICE_IN_PER_1K, price_out_per_1k: PRICE_OUT_PER_1K, cost: 0.0, status: ffailed:{type(last_error).__name__ if last_error else unknown}, latency_ms: int((time.monotonic() - start) * 1000), retry_count: MAX_RETRIES, rate_limit_wait_ms: rate_limit_wait_ms, trace_id: trace_id, } ) raise last_error or RuntimeError(review_failed) async def main(input_path: str, output_path: str) - None: init_db() with open(input_path, r, encodingutf-8) as f: items [json.loads(line) for line in f if line.strip()] global_sem asyncio.Semaphore(GLOBAL_CONCURRENCY) async with httpx.AsyncClient() as client: async def runner(item: dict[str, Any]) - dict[str, Any]: async with global_sem: return await review_one(client, item) results await asyncio.gather(*(runner(x) for x in items), return_exceptionsTrue) with open(output_path, w, encodingutf-8) as f: for result in results: if isinstance(result, Exception): f.write(json.dumps({error: str(result)}, ensure_asciiFalse) \n) else: f.write(json.dumps(result, ensure_asciiFalse) \n) if __name__ __main__: asyncio.run(main(reviews.jsonl, reviews_audited.jsonl))运行方式python batch_review.py如果某个模型不支持response_format可以删掉该字段但提示词里仍要强调“只输出 JSON”。解析失败建议在正式脚本里单独捕获json.JSONDecodeError并写入statusparse_error不要和模型审核失败混在一起。4. 限流规则全局令牌桶、租户滑动窗口与 429 退避并发参数只是第一层真正稳定的是限流规则。建议至少叠加四层全局令牌桶控制所有审核请求的总速率例如RATE_RPS8、RATE_BURST16。令牌桶允许短时突发但长期平均速率受控。租户滑动窗口每个租户每 60 秒最多 300 次审核防止单个团队刷满全局配额。并发信号量全局 16、租户 4避免大量长连接挂起。预算熔断每个租户每日成本达到 80% 时降并发达到 100% 时只处理priority0的实时任务批量任务排队到次日。滑动窗口限流器可以用下面这段代码实现# sliding_window.py import asyncio import collections import time class SlidingWindowLimiter: def __init__(self, max_requests: int, window_seconds: int): self.max_requests max_requests self.window_seconds window_seconds self.events: dict[str, collections.deque] collections.defaultdict(collections.deque) self.lock asyncio.Lock() async def acquire(self, tenant: str) - None: async with self.lock: now time.monotonic() q self.events[tenant] while q and now - q[0] self.window_seconds: q.popleft() if len(q) self.max_requests: sleep_for self.window_seconds - (now - q[0]) await asyncio.sleep(max(sleep_for, 0)) return await self.acquire(tenant) q.append(now)对 429 的处理不要只做无脑重试。优先读取Retry-After没有该头时使用指数退避0.5s、1s、2s并加入随机抖动。连续三次 429 后对该租户熔断 30 秒避免重试风暴拖垮其他任务。重试次数要写进成本对账表因为重试也会消耗 Token。下面是一张限流规则建议表规则键阈值超限动作全局令牌桶全局8 rps突发 16等待令牌租户滑动窗口tenant_id300 次 / 60s排队或降级全局并发全局16等待信号量租户并发tenant_id4等待信号量日成本预算tenant_id80% / 100%降并发 / 仅实时失败熔断tenant_id连续 3 次 429熔断 30s优先级队列也要有。priority0可以理解为交互式审核例如审核员正在页面等待结果priority1是批量历史回复priority2是失败补偿。不要让补偿任务和实时任务抢同一个信号量否则用户侧会感觉系统变慢。简单做法是开两个独立队列分别使用不同并发池但共用全局令牌桶。5. 成本对账表SQLite 表结构、写入字段与对账 SQL成本对账不是财务月底才做的事而是每次批处理运行时都要写。TaoToken 官网https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_cost控制台可以查看 Key 和用量但工程侧仍要在本地记录请求级明细。原因很简单控制台给的是总账你要能回答某个租户、某个模型、某次重试、某个审核批次分别花了多少。成本公式cost input_tokens / 1000 * price_in_per_1k output_tokens / 1000 * price_out_per_1k表结构已经在上一节脚本中创建。关键字段包括request_id、tenant_id、model、input_tokens、output_tokens、cost、status、retry_count、rate_limit_wait_ms。下面这段 SQL 只在本地 SQLite 执行不要把它接到生产库或审核系统的常驻连接上-- 在本地 SQLite 执行近 7 天租户/模型成本汇总 SELECT date(ts, unixepoch, localtime) AS day, tenant_id, model, SUM(input_tokens) AS input_tokens, SUM(output_tokens) AS output_tokens, ROUND(SUM(cost), 6) AS total_cost, SUM(CASE WHEN status ! ok THEN 1 ELSE 0 END) AS failed_requests, ROUND(AVG(latency_ms), 0) AS avg_latency_ms, ROUND(AVG(rate_limit_wait_ms), 0) AS avg_wait_ms FROM request_logs WHERE ts strftime(%s, now, -7 day) GROUP BY day, tenant_id, model ORDER BY day DESC, total_cost DESC;还可以按审核结果做质检对账-- 在本地 SQLite 执行按 verdict 和风险标签统计 SELECT json_extract(audit_json, $.verdict) AS verdict, COUNT(*) AS cnt, ROUND(AVG(json_extract(audit_json, $.relevance)), 3) AS avg_relevance, ROUND(AVG(json_extract(audit_json, $.template_score)), 3) AS avg_template, ROUND(AVG(json_extract(audit_json, $.sycophancy_score)), 3) AS avg_sycophancy FROM review_results GROUP BY verdict ORDER BY cnt DESC;如果你把模型返回结果单独存表建议不要只存audit_json还要把verdict、relevance、template_score、sycophancy_score抽成普通列。这样对账 SQL 不需要 JSON 函数查询更快。成本对账表可以按下面格式输出到 CSV日期租户模型输入 Token输出 Token成本失败数平均延迟平均限流等待2025-06-01team_aYOUR_MODEL_ID120000180000.08722380ms120ms2025-06-01team_bYOUR_MODEL_ID8000090000.053501760ms80ms一旦发现某租户avg_wait_ms持续升高说明限流参数过紧或租户任务量突增如果failed_requests高但total_cost低可能是 Key、模型 ID 或网络问题如果output_tokens异常大要检查审核提示词是否被回复正文带偏导致模型输出过长。6. Claude Code、Codex 与 CC Switch 三件套审核员交互与批处理配置分离批量审核脚本负责跑任务审核员日常排查则需要 Claude Code 或 Codex。这里要严格区分配置Claude Code 使用settings.json和ANTHROPIC_*环境变量Codex 使用config.toml不要把ANTHROPIC_*套到 Codex。CC Switch 三件套可以理解为三份配置文件Claude Code 配置、Codex 配置、共享环境变量。它们都指向同一个 TaoToken Key 和 Base URL但变量名和配置格式不同。Claude Code 的settings.json示例{ env: { ANTHROPIC_BASE_URL: https://taotoken.net/api, ANTHROPIC_API_KEY: YOUR_API_KEY, ANTHROPIC_MODEL: YOUR_MODEL_ID } }Claude Code 只认ANTHROPIC_*。ANTHROPIC_BASE_URL填https://taotoken.net/api不要带 UTM。模型 ID 以 TaoToken 控制台为准。Codex 的config.toml示例model YOUR_MODEL_ID model_provider taotoken [model_providers.taotoken] name TaoToken base_url https://taotoken.net/api env_key TAOTOKEN_API_KEY wire_api chatCodex 使用TAOTOKEN_API_KEY这种环境变量名不写ANTHROPIC_API_KEY也不要把ANTHROPIC_BASE_URL放进 TOML。共享环境变量文件可以这样写# .env.taotoken TAOTOKEN_API_KEYYOUR_API_KEY TAOTOKEN_BASE_URLhttps://taotoken.net/api TAOTOKEN_MODELYOUR_MODEL_IDCC Switch 三件套的推荐分工配置文件用途禁止事项Claude Codesettings.json审核员交互排查不要混用 Codex 的 TOMLCodexconfig.toml本地代码/文本分析不要写ANTHROPIC_*共享环境.env.taotoken批处理脚本与 CI不要提交真实 Key如果审核员需要查看某条回复为什么被判“模板腔”可以本地启动 Claude Code让它读取脱敏后的 JSONL 样本而不是直接连接线上数据库。批量任务里的 SQLite 对账表也只在本地查询。不要把 MCP、Agent 或脚本直连生产库审核数据要先脱敏导出再由本地任务读取。7. 上线检查与 CTA从模型对话到 Claude Code 文档上线前按下面清单过一遍Key 是否来自 TaoToken 控制台代码里是否只保留YOUR_API_KEY占位。Base URL 是否统一为https://taotoken.net/api工具配置是否误带 UTM。输入 JSONL 是否脱敏是否移除了手机号、邮箱、订单号等敏感字段。全局并发、租户并发、令牌桶、滑动窗口是否全部启用。429 是否有Retry-After优先、指数退避和租户熔断。每次请求是否写入request_logs成本公式是否使用真实价格。是否区分ok、http_429、parse_error、failed避免污染审核通过率。Claude Code 和 Codex 配置是否分离Codex 是否误用了ANTHROPIC_*。成本对账 SQL 是否只在本地 SQLite 执行是否避免了生产库直连。如果你还没有创建 Key可以先到 TaoToken 官网https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_cta完成注册和控制台配置。高转化路径建议按顺序走模型对话https://taotoken.net/models/detail/chat?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_chatCoding Planhttps://taotoken.net/coding-plan?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_plan创建 Keyhttps://taotoken.net/console/api-keys?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_keysClaude Code 文档https://taotoken.net/doc/ClaudeCodeAnthropic?utm_sourcetaotoken_aicg_blog_endutm_contentbatch_review_claudecode把YOUR_API_KEY替换掉把YOUR_MODEL_ID换成控制台当前可用模型先跑 50 条样本验证 JSON 解析、限流等待和成本写入再逐步放大并发。这样批量审核 ChatGPT 回复就不再是一次性脚本而是一条可观测、可限流、可对账的后端流水线。
返回列表