第11章:流式 API —— WebSocket 三种模式
约 3 分钟 · 更新于 2026-09-02
请求-响应之外,Encore 用三个变体覆盖实时通信:streamIn(上行)、streamOut(下行)、streamInOut(双向)。底层是 WebSocket,但握手、类型、鉴权、追踪全部沿用你已经学过的那套约定。
一、三种模式怎么选
| 模式 | 方向 | 典型场景 |
|---|
| api.streamIn | 客户端 → 服务端 | 分片上传、客户端埋点/传感器数据回传 |
| api.streamOut | 服务端 → 客户端 | 日志推送、进度通知、行情/航变推送 |
| api.streamInOut | 双向 | 聊天、协同编辑、实时客服 |
连接建立方式:客户端发起一个 HTTP 握手 (handshake) 请求 → 服务端接受 → 升级为 WebSocket → 双方按声明的消息类型收发。
二、streamIn:上行流
typescript
import { api } from "encore.dev/api";
// 握手参数:来自路径/查询串/请求头
interface Handshake {
filename: Query<string>;
}
// 客户端发来的每条消息
interface Message {
data: string;
done: boolean;
}
interface Response {
success: boolean;
size: number;
}
export const uploadStream = api.streamIn<Handshake, Message, Response>(
{ path: "/upload", expose: true },
async (handshake, stream) => {
const chunks: string[] = [];
for await (const msg of stream) { // 异步迭代收消息
chunks.push(msg.data);
if (msg.done) break;
}
await saveFile(handshake.filename, chunks.join(""));
return { success: true, size: chunks.length }; // 返回值作为最终响应发给客户端
},
);
类型参数 <Handshake, Message, Response>:握手数据、入站消息、最终响应。握手是可选的——不需要初始参数时写 api.streamIn<Message, Response>,handler 就只有 stream 一个参数。
三、streamOut:下行流
typescript
interface Handshake {
rows: Query<number>;
}
interface LogMessage {
row: string;
}
export const logStream = api.streamOut<Handshake, LogMessage>(
{ path: "/logs", expose: true },
async (handshake, stream) => {
for await (const line of tailLogs(handshake.rows)) {
await stream.send({ row: line }); // 主动推送
}
await stream.close(); // 服务端主动结束
},
);
四、streamInOut:双向流
typescript
interface InMessage { text: string; }
interface OutMessage { text: string; from: string; }
export const chatStream = api.streamInOut<InMessage, OutMessage>(
{ path: "/chat", expose: true },
async (stream) => { // 无握手版本:只有 stream 参数
for await (const msg of stream) {
await stream.send({ text: `echo: ${msg.text}`, from: "bot" });
}
},
);
五、广播模式:多客户端互通
聊天室、多人看板的标准写法——把活跃流存进 Map,收到消息遍历转发,断连清理:
typescript
import { api, StreamInOut } from "encore.dev/api";
const connected = new Map<string, StreamInOut<InMessage, OutMessage>>();
interface Handshake { userId: Query<string>; }
export const chat = api.streamInOut<Handshake, InMessage, OutMessage>(
{ path: "/chat", expose: true },
async (handshake, stream) => {
connected.set(handshake.userId, stream);
try {
for await (const msg of stream) {
// 广播给所有在线客户端
for (const [uid, peer] of connected) {
try {
await peer.send({ text: msg.text, from: handshake.userId });
} catch {
connected.delete(uid); // 发送失败视为断连
}
}
}
} finally {
connected.delete(handshake.userId); // 正常/异常退出都清理
}
},
);
内存态的边界
Map 存流对象是单实例内存态:水平扩容成多实例后,两个用户可能连在不同实例上,广播互相看不见。跨实例广播要把消息过一遍 Pub/Sub(第 7 章):收到消息 → publish → 每个实例订阅后推给本实例的连接。
六、客户端消费
生成的类型化客户端(第 15 章)对流式端点返回流对象:
typescript
const stream = client.chat.chatStream();
// 发送
await stream.send({ text: "hello" });
// 接收(异步迭代)
for await (const msg of stream) {
console.log(msg.from, msg.text);
}
// 底层 socket 事件
stream.socket.on("error", (e) => { /* 网络异常 */ });
stream.socket.on("close", (e) => { /* 断连处理,通常做退避重连 */ });
服务间也能互开流:
typescript
import { monitor } from "~encore/clients";
const stream = await monitor.logStream({ rows: 100 });
七、鉴权与追踪
- 鉴权:选项加 auth: true 即可,鉴权在握手阶段执行(走第 10 章的 auth handler),handler 内 getAuthData() 照常可用;
- 追踪:流式端点同样进分布式追踪,仪表盘可以看到连接生命周期。
八、常见误区
- 用轮询 HTTP 接口模拟推送——有 streamOut 就别让客户端每秒打一次接口。
- 广播 Map 忘了 finally 清理——断连的流残留导致 send 反复抛错、内存泄漏。
- 多实例部署后广播"丢消息"——单实例内存态问题,过 Pub/Sub 桥接。
- 握手想传参数却写在消息类型里——路径/查询串/头参数属于 Handshake 类型参数。
- streamIn 忘了 break 条件——客户端不关流时 for await 永远不结束,配合 done 标志或超时。
- 长连接里做重 CPU 计算——阻塞事件循环影响所有连接,重活丢给队列。
九、实践练习
- 写一个 streamOut 的"订单状态推送":握手带 orderNo,服务端每秒查一次状态有变化就推送,状态到终态后 close()。
- 把练习 1 的轮询查库改成订阅 Pub/Sub 事件触发推送(第 7 章 order-paid 主题复用)。
- 实现 streamIn 分片上传:客户端把文件切块发送,服务端拼接后存 Bucket(第 9 章),返回对象名。
- 实现最小聊天室(广播模式代码为骨架),两个终端连上互发消息;然后写出"扩到 3 个实例要改哪里"的方案。
- 给聊天室加 auth: true,把 from 改成从 getAuthData() 取用户名,消灭冒名顶替。
十、总结
- 三模式对应三方向:streamIn 上行、streamOut 下行、streamInOut 双向;底层 WebSocket,心智模型仍是类型化端点。
- Handshake 类型参数承载路径/查询串/头参数,可选;有它 handler 双参,没它单参。
- 收消息 for await,发消息 stream.send,服务端结束 stream.close。
- 广播 = Map 存流 + 遍历转发 + finally 清理;多实例场景必须过 Pub/Sub。
- 鉴权在握手期执行,auth: true + getAuthData() 与普通端点完全一致。
请继续阅读:第12章:中间件、CORS 与可观测性。
原始资料引用