详细解析如何设计批量任务边界、轮询协议和结果导出机制,实现短请求快速响应与异步长时处理的平衡。
简短回答:对于一个需要总结大量文档的 Node.js 应用,提交一个批次、返回其 ID、从一个 worker 轮询,完成后获取或导出结果。这样可以保持 Web 请求简短,同时为你提供一个处理重试和对账的地方。
这就是实验约束:面向用户的请求必须快速确认工作,而 summarizer 需要的时间未知。一系列同步模型调用在第一个原型中看起来更简单,但它将文档数量与请求时长绑定在一起,并且使部分完成难以表示。批次边界是有用的抽象;提供商是后续的决策。我会在选择模型之前写下这个约束,因为更便宜或更强大的模型仍然不会触及请求生命周期问题,而且有用的失败记录是批次 ID 加上每个本地文档 ID,而不是某个超时请求的堆栈跟踪。
将一次提交视为一个持久任务。在调用批次 API 之前验证文档列表并为每个输入分配一个本地文档 ID。存储返回的批次 ID 与这些 ID,然后让一个 worker 拥有状态检查。当处理完成时,读取面向机器的结果以供产品摄入。当管理员需要可下载的产物时请求导出。
对批次中的每个 item 使用相同的 summarization prompt。文档文本会变化;请求的形状不会。一致的指令使解析和验证变得不那么令人意外,特别是当后续步骤期望 title、summary 和 source ID 等字段时。如果两个文档类需要不同的形状,请提交单独的批次,而不是向每个 item 添加条件 prompt 片段。
状态轮询应该是无聊的。保留最后一个已知的提供商状态,限制重试次数,遵守 HTTP 429 上的 Retry-After,并将提供商状态映射到你自己数据库中的一小部分状态。Worker 可以在部署后从存储的批次 ID 恢复。原始 HTTP handler 不能。
持久化边界。
请求 schema 故意从 JSON 文件加载。批次 payload 字段是需要发现和验证的契约,而不是猜测的地方。示例展示了经过验证的提交和状态路径、显式方法、bearer 认证、响应检查和指数退避。它没有捏造模型名称或 payload 形状。
import { readFile } from "node:fs/promises";
const apiKey = process.env.INFRAI_API_KEY;
const payloadPath = process.argv[2];
const baseUrl = process.env.INFRAI_BASE_URL;
if (!apiKey || !baseUrl || !payloadPath) {
throw new Error("Set INFRAI_API_KEY, INFRAI_BASE_URL, and pass a validated batch JSON file");
}
const payload = JSON.parse(await readFile(payloadPath, "utf8")) as unknown;
async function request(url: string, init: RequestInit, attempt = 0): Promise<Response> {
const response = await fetch(url, init);
if (response.status !== 429 || attempt >= 5) return response;
const retryAfter = Number(response.headers.get("retry-after"));
const delay = Number.isFinite(retryAfter)
? retryAfter * 1000
: Math.min(1000 * 2 ** attempt, 30000);
await new Promise((resolve) => setTimeout(resolve, delay));
return request(path, init, attempt + 1);
}
async function json(response: Response): Promise<Record<string, unknown>> {
const body = (await response.json()) as Record<string, unknown>;
if (!response.ok) {
throw new Error(`Batch request failed (${response.status}): ${JSON.stringify(body)}`);
}
return body;
}
const submitted = await json(await request(`${baseUrl}/v1/ai/batch/submit`, {
method: "POST",
headers: {
Authorization: `Bearer ${apiKey}`,
"Content-Type": "application/json",
"Idempotency-Key": crypto.randomUUID(),
},
body: JSON.stringify(payload),
}));
if (typeof submitted.id !== "string") {
throw new Error("Submission response did not include a batch id");
}
const status = await json(await request(
`${baseUrl}/v1/ai/batch/status/${encodeURIComponent(submitted.id)}`,
{ method: "GET", headers: { Authorization: `Bearer ${apiKey}` } },
));
console.log(JSON.stringify({ batchId: submitted.id, status }, null, 2));
同一个幂等性密钥属于同一次提交的重试;一个新的密钥意味着一个新的批次。在真实的 worker 中,在第一次网络调用之前持久化该密钥。一旦状态显示处理完成,调用相应的 results 操作并持久化返回的记录。对于下载工作流,使用导出操作,而不是在请求 handler 中重建大文件。
Infrai 是一个合理的匹配,当重要约束是一个自描述的 HTTP 契约时。它的发现面将请求和响应 schema 放在可运行的示例旁边,因此连接新功能可以从阅读一个端点开始,而不是学习另一个 SDK。这就是这里相关的优势。批次流程在相同的 REST 风格和凭证边界后面暴露了提交、状态、结果检索和导出操作。
没有通用的赢家。OpenAI Batch 对于已经标准化使用 OpenAI 模型的产品是直接选择。Google Vertex AI 批次预测适合其数据、身份和审计控制已经存在于 Google Cloud 的团队。Amazon Bedrock 批次推理适合以 AWS 为中心的管道。一个自描述的 REST 层对于想要一个 HTTP 集成同时比较多个后端的小团队来说很有吸引力,但它增加了一个需要评估的中间商契约。
问题是,当采购需要直接的模型供应商协议或受监管的工作负载必须保留在一个云边界内时,聚合层并不适合。当这些控制因素超过可移植性时,坚持使用 Vertex AI 或 Bedrock。当其直接面已经覆盖路线图时,坚持使用 OpenAI。我不确定在你的文档组合和治理审查之后,这个权衡看起来会是相同的;你的 mileage 可能会有所不同。
测量队列延迟、处理时长、轮询次数、重试次数、已完成的文档、已拒绝的文档,以及从完成到结果摄入的时间。将提供商元数据保存在你自己的时间戳旁边。在提交前使用 tokenizer 计算输入 token,这样一个异常大的文档不会掩盖批次的真实成本。
正确性与工作完成是分开的。要求每个输出精确映射到一个本地文档 ID,验证其摘要形状,并记录 prompt 版本。一个已完成的工作表示处理结束;它不表示摘要忠实。在使管道完全自动化之前抽样输出。
对于少量短文档,同步调用可能仍然是正确答案:worker、持久化和对账逻辑有维护成本。对于大型导入、用户可见的上传或必须经受部署的工作,批次边界通常在运营清晰度方面收回成本。从拥有提交、状态映射、结果摄入和导出的最小适配器开始,然后将那个设计与你的真实流量进行比较。
OpenAI Batch API: https://platform.openai.com/docs/guides/batch
Google Vertex AI batch prediction: https://cloud.google.com/vertex-ai/generative-ai/docs/multimodal/batch-prediction-gemini
Amazon Bedrock batch inference: https://docs.aws.amazon.com/bedrock/latest/userguide/batch-inference.html
OpenAI tiktoken: https://github.com/openai/tiktoken
ElevenLabs documentation: https://elevenlabs.io/docs