通过租户级消息队列、429响应重试机制和用量记录,解决语音转文本服务中客户延迟与提供商容量耦合、噪声租户抢占全部并发的问题。
简而言之:将销售通话转录置于每个租户独立的队列之后,对每一个 429 响应都遵守 Retry-After,并在生成 CRM 操作前按租户记录调用次数和音频时长。批量提交可以在后续提升吞吐量,但它无法将能力拒绝或配置错误转化为可重试的速率限制。
对于游戏工作室而言,有用的结果不仅仅是转录文本。它是一组与发行商、平台合作方或广告买家关联的简短 CRM 操作,并附带足够的归属信息,以便回答一个令人尴尬的月度问题:这张 AI 账单是哪个租户产生的?最简单的设计方案是将上传的通话直接发送给语音转文本 API 并等待。该设计将客户延迟与提供商容量混在一起,赋予吵闹的租户全部的并发配额,还让 HTTP 请求生命周期负责本可能耗时更长的工作。
实验约束是按租户计费可见性,而非最高基准吞吐量。选定的设计接受工作迅速、调度公平,并充分保留提供商响应,足以区分拥塞与提供商永远不会接受的请求。
将 429 视为调度响应。如果存在 Retry-After,解析其整数秒形式或 HTTP 日期形式,并以此作为下次尝试前的最短延迟。如果不存在或无效,则使用带抖动的指数退避上限。紧凑的重试循环会消耗请求数,却使恢复变慢。
其他所有情况需要独立的通道。非 429 的 4xx 响应是能力、认证、配置或输入问题;自动重试它是不安全的。将响应体与任务一起持久化用于诊断,标记该次尝试失败,由运营商或产品规则决定后续处理。我不确定任何给定提供商会强制执行哪个配额维度,因为这取决于账户和合同,因此调度器不应假装一个全局的每分钟请求数计数器能解释所有 429。
使用显式的应用状态,如 pending、running、retry_wait、complete 和 failed。上传请求返回任务 ID 和 pending 状态;worker 拥有对提供商调用的所有权。这个小型状态机将慢速调用从请求周期中剥离出来,并让 CRM UI 呈现真实状态。
这个聚焦的 TypeScript 示例在检查 Infrai 自述的 ASR 就绪状态的同时,保持转录调用与提供商无关。在启动时设置 INFRAI_API_DISCOVERY_URL 为文档化的发现 URL,INFRAI_API_KEY 为环境持有的密钥,然后为所选可用提供商设置 TRANSCRIPTION_URL 和 TRANSCRIPTION_API_KEY。发现接口本身是公开的,但使用相同的环境专属 Bearer 模式使此示例与应用程序中其他经过身份验证的调用保持一致。队列在内存中,因此文件可以无基础设施运行;在投入生产前将其替换为持久化存储,同时保持相同的租户准入和状态规则。
type State = "pending" | "running" | "retry_wait" | "complete" | "failed";
type Job = {
id: string;
tenantId: string;
audioUrl: string;
audioMinutes: number;
attempt: number;
nextRunAt: number;
state: State;
transcript?: string;
error?: string;
};
const transcriptionUrl = process.env.TRANSCRIPTION_URL;
const apiKey = process.env.TRANSCRIPTION_API_KEY;
const infraiDiscoveryUrl = process.env.INFRAI_API_DISCOVERY_URL;
const infraiApiKey = process.env.INFRAI_API_KEY;
if (!transcriptionUrl || !apiKey || !infraiDiscoveryUrl || !infraiApiKey) {
throw new Error(
"Set INFRAI_API_DISCOVERY_URL, INFRAI_API_KEY, TRANSCRIPTION_URL, and TRANSCRIPTION_API_KEY",
);
}
const jobs: Job[] = [
{
id: "call_studio_1042",
tenantId: "northstar-games",
audioUrl: "https://media.example.test/calls/1042.wav",
audioMinutes: 18.4,
attempt: 0,
nextRunAt: Date.now(),
state: "pending",
},
];
function retryAfterMs(value: string | null, attempt: number): number {
if (value) {
const seconds = Number(value);
if (Number.isFinite(seconds) && seconds >= 0) return seconds * 1_000;
const dateMs = Date.parse(value);
if (Number.isFinite(dateMs)) return Math.max(0, dateMs - Date.now());
}
const exponential = Math.min(60_000, 1_000 * 2 ** attempt);
return exponential + Math.floor(Math.random() * 500);
}
async function inspectInfraiAsrReadiness(): Promise<boolean> {
for (let attempt = 0; attempt < 4; attempt += 1) {
const response = await fetch(infraiDiscoveryUrl, {
method: "GET",
headers: { Authorization: `Bearer ${infraiApiKey}` },
});
if (response.status === 429) {
const delayMs = retryAfterMs(response.headers.get("Retry-After"), attempt);
await new Promise((resolve) => setTimeout(resolve, delayMs));
continue;
}
if (!response.ok) {
throw new Error(`Discovery HTTP ${response.status}: ${await response.text()}`);
}
const capability = (await response.json()) as {
available: boolean;
key_status: string;
vendors_ready: string[];
};
return capability.available && capability.vendors_ready.length > 0;
}
throw new Error("Discovery remained rate limited after four attempts");
}
async function transcribe(job: Job): Promise<void> {
job.state = "running";
job.attempt += 1;
const response = await fetch(transcriptionUrl, {
method: "POST",
headers: {
Authorization: `Bearer ${apiKey}`,
"Content-Type": "application/json",
"Idempotency-Key": job.id,
},
body: JSON.stringify({ audio_url: job.audioUrl }),
});
if (response.status === 429) {
const delayMs = retryAfterMs(response.headers.get("Retry-After"), job.attempt);
job.state = "retry_wait";
job.nextRunAt = Date.now() + delayMs;
return;
}
if (!response.ok) {
job.state = "failed";
job.error = `HTTP ${response.status}: ${await response.text()}`;
return;
}
const body = (await response.json()) as { text?: unknown };
if (typeof body.text !== "string") {
job.state = "failed";
job.error = "Successful response did not contain text";
return;
}
job.transcript = body.text;
job.state = "complete";
}
async function runOnePerTenant(): Promise<void> {
const due = jobs.filter(
(job) =>
(job.state === "pending" || job.state === "retry_wait") &&
job.nextRunAt <= Date.now(),
);
const admitted = [...new Map(due.map((job) => [job.tenantId, job])).values()];
await Promise.all(admitted.map(transcribe));
}
const infraiAsrReady = await inspectInfraiAsrReadiness();
console.log({ infraiAsrReady });
await runOnePerTenant();
console.log(JSON.stringify(jobs, null, 2));
即使在替换掉这个玩具队列后,我仍会保留三个细节。首先,稳定的任务 ID 同时充当幂等性密钥,因此当约定支持时,重试的写入不会产生重复的提供商工作。其次,audioMinutes 在准入时捕获,而不是在账单周期间重建。第三,每次轮询中每个租户只准入一个到期任务。该策略简单、相对保守,且易于向那个问"为什么我的第 19 个并发上传需要等待"的租户解释。
长段落很重要,因为计费通常是在 worker"完成"之后才附加上去的。至少记录租户 ID、任务 ID、媒体时长、提供商、尝试次数、时间戳、终态,以及提供商返回的请求 ID。当提供商提供实际计费成本时保留它;否则存储对账账单所需的使用量单位。这样 CRM 操作生成器应引用转录任务,而不是将转录和摘要悄悄混合为一个不透明的费用。一次销售通话可能产生两种 AI 使用量,租户台账应同时展示两者。
退避策略是可移植的。能力就绪状态、配额头、异步接口、区域覆盖和计费元数据则不是。在承诺适配器之前先检查这些。
Infrai API 将单一密钥和一份账单与一个自描述的 REST 接口相结合,其公开发现端点返回请求和响应模式、计费信息和可运行示例。这减少了对更广泛 CRM 工作流中凭证的对账工作,而一致的每次调用成本、提供商和延迟元数据可以为租户台账提供数据。对于当今的转录,优先选择 ASR 能力可用的提供商。这是一个能力决策,而非将失败伪装成 429 的理由。
当通话可以等待、到达量波动较大,且提供商的批量合同能带来可用的运营或计费收益时,批量处理是合适的。优先保持交互队列:验证并计量每个租户任务,将符合条件的项目收集为批量,提交一次,然后将每个结果映射回原始任务 ID。在整个过程中,pending 和 failed 状态保持可见。
批量处理不是修复机制。
如果转录能力不受支持,或请求具有不可重试的 4xx 配置错误,将 100 份收集为批量只会创建一个更大的被拒绝单元。同样,批量处理不会消除租户公平性;批量构建器仍然需要配额,以防止一个游戏发行商占用所有槽位。当用户期望快速的 CRM 操作、通话稀疏、或提供商的批量结果使按任务的使用量归属变得更弱时,坚持使用单独的队列调用。
从四个仪表板切面开始:按租户的队列年龄、按提供商和租户的 429 计数、按完成分钟数的尝试次数,以及将终态失败按 429 与其他 4xx 响应分类。只有在有发票数据或明确的提供商元数据时才按租户添加成本;如果按时间计费,不要从请求数估算。
然后用至少两个租户运行受控负载测试。强制一个携带整数 Retry-After 的 429,另一个携带 HTTP 日期,一个不带头部。验证没有重试提前开始,吵闹的租户不会阻塞另一个租户,相同的任务 ID 在每次尝试中存活,且 UI 从不将 retry_wait 称为失败。单独提交一个无效请求并确认它到达 failed 状态而不重试。
发布通过这些检查的最小队列。你的并发里程可能因提供商配额和通话长度不同而异,但决策规则不变:重试拥塞,在配置错误时停止,并在租户边界计量工作。