详解API Gateway WebSocket架构下连接管理API(@connections)的调用机制,说明为何流式处理需要额外的连接ID存储和排序逻辑,对比SSE/chunked HTTP的区别。
连接与 token 生成不在同一处
在 Server-Sent Events 或分块 HTTP 响应中,一个进程拥有响应流并在 token 到达时写入。而 API Gateway WebSocket 中,API Gateway 拥有 socket。你的集成在每条入站消息时调用、返回、然后消亡。后端想要发送的任何内容都通过 @connections 管理 API 发出——一个普通的、带签名的请求/响应 HTTPS 调用,通过 connection id 来寻址。
这一设计所有别扭之处都源于一个事实:connection id 必须存在某个持久化的地方,这样另一个进程才能找到它。顺序是你自己的问题,因为对同一连接的两个并发 POST 就是两个独立的 HTTP 请求。断连是通过调用失败来发现的,而不是 socket 在你手下关闭。
路由通过 routeSelectionExpression 完成,这是 API 本身设置的一个属性,指定入站消息中的一个 JSON 属性名——通常是 "$request.body.action"。三个路由 key 是预定义的:$connect 在连接建立时调用,$disconnect 在任意一方离开时调用,$default 捕获不匹配任何路由或无法被解析为 JSON 的消息。
$connect 与 connection id
在 $connect 时授权,而不是之后。 这是生命周期中唯一拥有正常请求上下文可供授权的时机——查询字符串、headers、authorizer——从它返回非 2xx 状态码会直接拒绝连接。之后,每一帧你拥有的只有 connection id。
export const handler = async (event) => {
const { connectionId, requestTimeEpoch } = event.requestContext;
const userId = event.requestContext.authorizer?.principalId;
if (!userId) return { statusCode: 401 };
await ddb.send(new PutItemCommand({
TableName: "ws-connections",
Item: {
connectionId: { S: connectionId },
userId: { S: userId },
// 2h max lifetime plus slack; DynamoDB TTL sweeps the leftovers
expiresAt: { N: String(Math.floor(requestTimeEpoch / 1000) + 8000) },
},
}));
return { statusCode: 200 };
};
TTL 属性不是可选的琐事。因为连接可能在 $disconnect 从未运行过的情况下结束——网络丢失、突发关闭——表里会积累命名已不存在 socket 的行,而每一个都是一次未来必然失败的 API 调用。设置一个远超两小时最大连接生命周期的 TTL,可以限制这种增长,而无需你写一个清理程序。
用 @connections 推送 token
生成 token 的进程——第二个 Lambda、Step Functions 任务、消费模型流的容器——通过管理端点 POST 每一个 chunk,管理端点即 API 的 execute-api 主机名加上 stage 名称:
import {
ApiGatewayManagementApiClient,
PostToConnectionCommand,
} from "@aws-sdk/client-apigatewaymanagementapi";
const mgmt = new ApiGatewayManagementApiClient({
endpoint: "https://abc123.execute-api.us-east-1.amazonaws.com/prod",
});
async function relay(connectionId, stream) {
let buffer = "";
let lastFlush = Date.now();
for await (const chunk of stream) {
buffer += chunk.delta ?? "";
// One POST per token is one signed HTTPS request per token.
if (buffer.length > 200 || Date.now() - lastFlush > 100) {
await mgmt.send(new PostToConnectionCommand({
ConnectionId: connectionId,
Data: JSON.stringify({ type: "delta", text: buffer }),
}));
buffer = "";
lastFlush = Date.now();
}
}
if (buffer) {
await mgmt.send(new PostToConnectionCommand({
ConnectionId: connectionId,
Data: JSON.stringify({ type: "delta", text: buffer }),
}));
}
await mgmt.send(new PostToConnectionCommand({
ConnectionId: connectionId,
Data: JSON.stringify({ type: "done" }),
}));
}
缓冲是值得保留的部分。一个模型以每秒 60 个 token 的速度输出到一个无缓冲的 relay,每秒每个读者就是 60 个带签名的 HTTPS 请求,每一个都会被计费,并且都与你的普通流量共享同一个账户级 API Gateway 节流配额——参见 AWS API Gateway 的节流限制如何相互作用。按字符阈值或 100ms 计时器刷新,保持界面响应的同时将调用次数降低一个数量级。
调用方的 IAM 角色需要 execute-api:ManageConnections,这是一个与控制 API 调用的 execute-api:Invoke 不同的 action:
{
"Effect": "Allow",
"Action": "execute-api:ManageConnections",
"Resource": "arn:aws:execute-api:us-east-1:111122223333:abc123/prod/POST/@connections/*"
}
帧、消息和时长限制
AWS 文档中记录了四个数字,它们塑造了这一设计,截至撰写本文时这四个数字均不可调整:WebSocket 帧大小 32 KB,消息负载 128 KB,最大连接时长 2 小时,空闲连接超时 10 分钟。超过帧大小的消息必须拆分到多个帧。
API Gateway 在关闭时返回的状态码告诉你撞上了哪个限制,在客户端中值得显式处理。码 1001 涵盖了 10 分钟空闲超时和 2 小时生命期上限——同一代码对应两种截然不同的情况,所以客户端应该重连而不是诊断。1009 表示消息太大无法处理。1003 表示收到了二进制数据;WebSocket API 不支持二进制媒体类型,这就排除了通过此传输发送原始音频帧的可能。1008 在客户端发送请求过多时返回。
这些配额以及每秒 500 个新连接的账户限制来自撰写本文时的 Amazon API Gateway 配额页面。连接速率限制是可调整的;帧、消息和时长限制文档中标注为不可调整。
10 分钟空闲超时是最令人意外的一个,因为"空闲"意味着双向都没有流量。用户在阅读一段长回复时什么都不发送,一旦生成结束也什么都不收,所以聊天会话会在读者注意力正集中的时候掉线。客户端每隔几分钟向一个自定义路由发送 ping 即可重置它。2 小时上限无法被任何操作重置;设计客户端重连并重新认证,而不是把一个 socket 当作一个会话。
GoneException 与清理问题
当你向一个已关闭的连接 POST 时,管理 API 返回 HTTP 410,SDK 抛出 GoneException。这是了解断连的正常方式,不是错误状况,应该在调用处处理——删除存储的行并放弃流:
try {
await mgmt.send(new PostToConnectionCommand({ ConnectionId, Data }));
} catch (err) {
if (err.name === "GoneException") {
await ddb.send(new DeleteItemCommand({
TableName: "ws-connections",
Key: { connectionId: { S: ConnectionId } },
}));
throw new ClientGoneError(); // stop consuming the model stream
}
throw err;
}
重新抛出而非吞掉异常,其重要性超乎表面。如果 relay 在读者已经离开后继续消费模型的流,你正在为永远没人会看到的输出 token 付费,而且是生成的全过程。收到 GoneException 时取消上游请求,断连的代价从完整生成变成零——在不稳定的移动网络下,断连并不罕见。
上面的 relay 假设了一种流形态。事实并非如此:提供商在事件帧结构、delta 中文本的位置、完成信号方式以及中流错误报告方式上都不同——所以针对一个提供商写的 relay,在加入备用方案那天就需要第二个解析器。Multigrid 将跨提供商的流式传输规范化为一种事件形态,这意味着 WebSocket 这端始终是单一代码路径,而背后的模型可以随时更换。