Agent X-Ray
RuntimeNotesAbout
Notes/代码工程/Encore/第11章

第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() 照常可用;
  • 追踪:流式端点同样进分布式追踪,仪表盘可以看到连接生命周期。

八、常见误区

  1. 用轮询 HTTP 接口模拟推送——有 streamOut 就别让客户端每秒打一次接口。
  2. 广播 Map 忘了 finally 清理——断连的流残留导致 send 反复抛错、内存泄漏。
  3. 多实例部署后广播"丢消息"——单实例内存态问题,过 Pub/Sub 桥接。
  4. 握手想传参数却写在消息类型里——路径/查询串/头参数属于 Handshake 类型参数。
  5. streamIn 忘了 break 条件——客户端不关流时 for await 永远不结束,配合 done 标志或超时。
  6. 长连接里做重 CPU 计算——阻塞事件循环影响所有连接,重活丢给队列。

九、实践练习

  1. 写一个 streamOut 的"订单状态推送":握手带 orderNo,服务端每秒查一次状态有变化就推送,状态到终态后 close()
  2. 把练习 1 的轮询查库改成订阅 Pub/Sub 事件触发推送(第 7 章 order-paid 主题复用)。
  3. 实现 streamIn 分片上传:客户端把文件切块发送,服务端拼接后存 Bucket(第 9 章),返回对象名。
  4. 实现最小聊天室(广播模式代码为骨架),两个终端连上互发消息;然后写出"扩到 3 个实例要改哪里"的方案。
  5. 给聊天室加 auth: true,把 from 改成从 getAuthData() 取用户名,消灭冒名顶替。

十、总结

  1. 三模式对应三方向:streamIn 上行、streamOut 下行、streamInOut 双向;底层 WebSocket,心智模型仍是类型化端点。
  2. Handshake 类型参数承载路径/查询串/头参数,可选;有它 handler 双参,没它单参。
  3. 收消息 for await,发消息 stream.send,服务端结束 stream.close
  4. 广播 = Map 存流 + 遍历转发 + finally 清理;多实例场景必须过 Pub/Sub。
  5. 鉴权在握手期执行,auth: true + getAuthData() 与普通端点完全一致。

请继续阅读:第12章:中间件、CORS 与可观测性


原始资料引用



本章目录
一、三种模式怎么选二、streamIn:上行流三、streamOut:下行流四、streamInOut:双向流五、广播模式:多客户端互通六、客户端消费七、鉴权与追踪八、常见误区九、实践练习十、总结原始资料引用Related Documents
苏ICP备2025204887号-2