详细对比了同步请求、有限并行和持久化异步任务三种边界场景,给出了何时用、何时换的具体决策规则和CRM导出验证模式。
对于市场商城的销售通话,当多份文档需要成为一组统一的 CRM 操作审核集时,使用异步任务;仅当一份简短文档能够在调用者的延迟预算内完成时,才使用内联请求。每个输入保留一个结果,暴露部分进度,且仅导出包含源 ID、结果和 schema 版本的记录。
从以下决策表开始:
关键的边界不在于"批处理与否",而在于所有权。如果 API 接收一个集合,服务应该通过最终结果和可验证的导出拥有该集合。不要让客户端从那些恰好已解析的 Promise 中重建真相。
Node.js 批处理摘要 API 在处理多份文档时应该做什么?
它应该将一个 admission 请求转换为一个稳定的任务记录,在声明的并发限制下处理每份文档,并在宣布任务完成前发布每项的结果。结果模型至少需要四个标识:任务、输入文档、处理尝试和导出。没有它们,重复提交可能看起来像新工作,重试可能覆盖有用的证据,而导出可能悄悄省略失败的通话。
以市场商城为例,假设一个卖家关于同一个账户有三次通话:需求发现、报价和法律审查。期望的 CRM 更新不仅仅是三段文字。它可能包括一个合并的后续步骤、一个负责人、一个异议,以及指向某个 transcript 的证据。并行运行三个 Promise 很容易。但判断法律通话的失败是否允许部分 CRM 更新,才是实际的产品决策。
使用显式状态机:accepted 表示请求和幂等 key 已存储;running 表示至少有一项可能在进行中;completed 表示每项都有最终结果;cancelled 表示不会启动新的尝试。如果调用者可能将其与完全失败混淆,就不要把 completed_with_errors 放在顶层状态。更清晰的契约是 completed 加上 succeeded、failed 和 skipped 项的计数。
用文字描述 diagram:admission 写入任务,队列释放 item ID,worker 获取文本并摘要,验证器检查结构化操作,reducer 决定是否允许跨通话合并,导出器冻结清单。指标观察每个箭头。数据库仍然是任务真相的来源。
这种分离很重要。队列回执只证明已投递,而非可用的摘要。
异步 Node.js 任务如何安全地导出批处理摘要结果?
在选择队列之前先定义契约。以下 TypeScript 类型使部分完成可见,并使每个导出的操作都可追溯到其文档。它们还将模型特定的响应形状保留在应用其余部分之外。
type JobState = "accepted" | "running" | "completed" | "cancelled";
type ItemState = "pending" | "running" | "succeeded" | "failed" | "skipped";
type DocumentInput = {
documentId: string;
accountId: string;
transcript: string;
};
type CrmAction = {
kind: "follow_up" | "update_stage" | "record_objection";
owner: string | null;
text: string;
sourceDocumentIds: string[];
};
type ItemResult = {
documentId: string;
state: ItemState;
attemptCount: number;
summary?: string;
actions?: CrmAction[];
errorCode?: "input_too_large" | "invalid_output" | "rate_limited";
};
type BatchJob = {
jobId: string;
idempotencyKey: string;
state: JobState;
schemaVersion: 1;
createdAt: string;
items: ItemResult[];
};
API 应该确认 admission,返回一个 job ID,并让调用者在不保持原始连接打开的情况下读取任务状态。单独的操作应该只读取一个最终快照。如果在 worker 仍在更新行时生成导出,同一个任务的两次下载可能不一致,即使两者看起来都有效。
幂等性填补了另一个缺口。将客户端的 key 绑定到有序输入 ID 和处理选项的规范摘要。带有相同摘要的重复 key 返回现有任务;相同的 key 但不同的摘要则产生冲突。这条规则防止网络重试生成重复的 CRM 操作,同时仍然暴露一个意外将 key 重用于不同工作的调用者。
考虑这个尴尬场景,因为这是普通 Promise.all 设计开始变得缺乏说服力的地方。客户端提交了需求发现、报价和法律审查的 transcript;需求发现最先完成,报价在速率限制后重试,而法律审查验证失败。在重试等待期间,客户端失去连接并再次提交相同请求。幂等记录应该将第二次请求指向原始任务,而不是再创建三项。一旦报价成功且法律审查达到重试限制,任务变为最终状态,包含两个成功和一个失败。其导出仍然包含三行,包括失败的法律行及其错误码,因此 CRM 导入器可以提出支持的操作,而不会假装集合已完成。如果策略要求三次通话都是强制性的,reducer 将每项提议的操作标记为仅供审查。相同的机制。不同的业务规则。
缺失的行是契约中的 bug。
以下是更深入的实现路径。它使用通用 summarizer 而不是供应商 SDK,限制并发,记录每个结果,并在所有项都达到最终状态之前拒绝生成导出。生产存储必须使状态转换原子化;接口使该需求保持可见。
type SummaryOutput = { summary: string; actions: CrmAction[] };
interface Summarizer {
summarize(input: DocumentInput): Promise<SummaryOutput>;
}
interface JobStore {
load(jobId: string): Promise<BatchJob>;
markRunning(jobId: string, documentId: string): Promise<void>;
saveSuccess(jobId: string, documentId: string, output: SummaryOutput): Promise<void>;
saveFailure(jobId: string, documentId: string, errorCode: ItemResult["errorCode"]): Promise<void>;
finishIfTerminal(jobId: string): Promise<void>;
}
async function runWithLimit<T>(
values: T[],
limit: number,
work: (value: T) => Promise<void>,
): Promise<void> {
const pending = [...values];
const workers = Array.from({ length: Math.min(limit, pending.length) }, async () => {
while (pending.length > 0) {
const value = pending.shift();
if (value !== undefined) await work(value);
}
});
await Promise.all(workers);
}
async function processJob(
jobId: string,
inputs: DocumentInput[],
summarizer: Summarizer,
store: JobStore,
): Promise<void> {
await runWithLimit(inputs, 4, async (input) => {
await store.markRunning(jobId, input.documentId);
try {
const output = await summarizer.summarize(input);
await store.saveSuccess(jobId, input.documentId, output);
} catch (error) {
const code = error instanceof RangeError ? "input_too_large" : "invalid_output";
await store.saveFailure(jobId, input.documentId, code);
}
});
await store.finishIfTerminal(jobId);
}
async function exportJob(jobId: string, store: JobStore): Promise<string> {
const job = await store.load(jobId);
if (job.state !== "completed") throw new Error("job_not_terminal");
return job.items
.map((item) => JSON.stringify({ jobId, schemaVersion: 1, ...item }))
.join("\n");
}
4 的限制是一个示例,不是通用调优值。衡量它。增加并发直到队列年龄改善,同时不将速率限制响应、内存压力或下游延迟推至超出预算。然后保留一些余量;在清洁测试流量下发现的设置对于早高峰时段的突发通话来说过于乐观。
Token 计数属于 admission 之前或项进入昂贵处理阶段之前。BPE tokenizer 可以估算模型输入大小,但编码细节取决于模型和 tokenizer 配置。对于硬限制,我不认为字符计数捷径值得这种歧义;请用所选运行时使用的 tokenizer 进行验证。tiktoken 项目文档记录了其 BPE 实现和支持的用法。
将质量和延迟作为同一个系统来观察
异步设计可能隐藏缓慢。通过分别测量队列延迟、处理时间和导出延迟来修复。单一本机百分比无法告诉你 worker 是慢了、容量不足,还是导出器在等待某个中毒的项。
跟踪任务 admission 数量、可运行队列年龄、项尝试次数、项持续时间、最终结果、验证失败和导出生成持续时间。在日志中附加 jobId 和 documentId,但将 transcript 文本和生成的摘要排除在常规遥测之外。高基数 ID 对日志和追踪有用;作为指标标签可能很昂贵或不可用,因此按有界维度(如结果和工作负载类别)聚合指标。
质量也需要操作信号。对于 CRM 操作,验证 schema、允许的操作类型、源文档引用和账户一致性。然后使用稳定 rubric 对已完成任务进行抽样人工审查:事实支持、遗漏的承诺、错误的负责人和不安全的状态变更。自动验证捕获格式错误的输出。它不能证明摘要保留了法律通话中的决定性句子。
设置两个服务目标,而不是将虚假妥协强制为一个数字:符合条件任务的目标完成率和审查结果的目标质量接受率。在目标完成率被违反之前就队列年龄持续增长发出警报。当 schema 有效输出失去人工接受度时单独发出警报,因为增加 worker 不会修复这种回归。
还有一个实际细节:重试应遵循错误类别。速率限制可以用有界退避和抖动重试;无效输入应立即终止;无效的结构化输出可能值得有限的重新生成尝试。存储尝试次数和最终代码。否则,消耗了五次尝试的任务看起来与干净成功的任务相同,容量规划变成了猜谜。
知道这个模式何时是错误的选择
问题是增加了机械结构。当用户需要在交互式编辑期间获取一个小摘要并且可以安全地重试整个请求时,持久化任务并不合适。保留内联路径。对于具有很小、可丢弃输入集且没有共享导出的内部脚本,有界的客户端并行也可能合理。
当每次通话属于不同账户、同意边界或保留类别时,不要使用自动跨文档合并。先拆分任务。同样,当策略要求人工批准时,不要直接从生成的文本发布 CRM 状态变更;导出带有证据的提议操作,让审查系统拥有决策权。
没有通用的并发限制、块大小或轮询间隔。transcript 长度、运行时配额、存储争用以及通话结束与 CRM 更新之间的可接受延迟都会影响你的结果。负载测试代表性分布,包括在一个长 transcript 夹杂几个短 transcript 中,并验证长项不能擦除或无限期延迟其周围已完成的结果。
先交付账本。优化可以随后进行。