rifts.to
← 所有文章

我如何用 Cloudflare Durable Objects 构建实时观众问卷

2026-04-04

我做了 rifts.to——一个给主讲人用的现场观众问卷工具。扫个二维码,回答一个问题,看着结果实时更新。表面上很简单。有意思的是底下真正让它跑起来的那部分,因为“边缘侧的实时扇出”是个听着容易、其实不容易的问题。

这篇讲的是 Cloudflare Durable Objects:它们是什么、为什么它们正好是解决这个具体问题的工具,以及实现出来到底长什么样。

问题:边缘侧的实时,比听上去难

场景是这样的。主讲人建了一份问卷,把二维码放出来。房间里一百个人扫码,开始提交答案。主讲人的面板需要实时更新——每一次提交都应该在一两秒内出现,而且不能靠轮询。

服务器推送事件(SSE)天然适合这件事。面板对服务器打开一条持久的 HTTP 连接,只要有变化,服务器就顺着这条连接把事件推下去。简单、单向、支持广泛,也没有 WebSocket 握手那些复杂度。

但有个坑:SSE 需要一条指向单个服务器进程的、持久且有状态的连接。当一份答卷落在某个边缘节点上时,你需要把事件推给那些可能开在其他边缘节点上的面板连接。在传统的无状态架构里,你会用一层发布订阅来解决——Redis、Kafka 之类的东西。

在边缘侧,这就别扭了。Cloudflare Workers 在设计上就是无状态的。你没法跨请求持有连接,也没法在多次调用之间共享内存。正是这一点让它们又快又便宜——也正是这一点让实时扇出变得困难。

于是轮到 Durable Objects 出场。

Durable Objects 究竟是什么

Durable Objects 是 Cloudflare 对“有状态的边缘计算”给出的答案。一个 Durable Object 是一个单实例的 JavaScript 类,它:

  • 只存在于一个位置,在 Cloudflare 的网络里(不跨区域复制)
  • 拥有持久化存储,通过内建的键值存储
  • 串行处理请求,在单个实例内——不存在并发问题
  • 可以按需长时间持有 WebSocket 或 HTTP 连接

关键在于那个单实例保证。当你按 ID 把流量路由到某个 Durable Object 时,带着这个 ID 的每一个请求都会打到同一个实例上。这正是做扇出所需要的原语:属于同一份问卷的所有 SSE 连接都住在同一个 Durable Object 里,所以新答卷到达时,它可以从内存里的同一个地方推给所有连接。

在 rifts.to 里,每一份问卷都有自己的 Durable Object——一个 SurveyRoom。这个房间持有该问卷结果面板的所有打开着的 SSE 连接。

架构

一次提交在系统里是这样流动的:

观众提交答案
        ↓
Cloudflare Pages (Next.js Edge Route: /api/respond)
        ↓
校验 Turnstile 并写入 D1(SQLite)
        ↓
按问卷 ID 取到 SurveyRoom Durable Object
        ↓
调用房间的 /broadcast
        ↓
SurveyRoom 把事件推给所有打开的 SSE 连接
        ↓
主讲人的面板即刻更新

面板会对 /api/admin/[token]/stream 打开一条连接——这是一个边缘路由,它用管理令牌做认证、查出问卷,然后把这条 SSE 连接直接交给 SurveyRoom DO。DO 让这条连接保持活着,并在每次 /broadcast 被调用时往里写数据。

注意这条流是按管理令牌而不是问卷 ID 来标识的。这意味着只有通过认证的主讲人面板才拿得到实时流;观众只是提交,然后看到一个致谢页面。

SurveyRoom 这个 Durable Object

这个 DO 暴露三个端点:/connect(订阅一个 SSE 客户端)、/broadcast(推给所有客户端)和 /close(发一个终止事件并关掉所有连接)。完整实现如下:

import type { SSEEvent } from "../lib/types";

export class SurveyRoom implements DurableObject {
  private connections: Set<WritableStreamDefaultWriter<Uint8Array>> = new Set();
  private encoder = new TextEncoder();

  constructor(
    private readonly state: DurableObjectState,
    private readonly env: CloudflareEnv
  ) {}

  async fetch(request: Request): Promise<Response> {
    const url = new URL(request.url);

    if (url.pathname === "/connect") {
      return this.handleConnect(request);
    }
    if (url.pathname === "/broadcast" && request.method === "POST") {
      return this.handleBroadcast(request);
    }
    if (url.pathname === "/close" && request.method === "POST") {
      return this.handleClose();
    }

    return new Response("Not found", { status: 404 });
  }

  private handleConnect(request: Request): Response {
    const { readable, writable } = new TransformStream<Uint8Array, Uint8Array>();
    const writer = writable.getWriter();
    this.connections.add(writer);

    const connectedEvent: SSEEvent = { type: "connected" };
    writer.write(this.encoder.encode(`data: ${JSON.stringify(connectedEvent)}\n\n`));

    const cleanup = () => {
      this.connections.delete(writer);
      writer.close().catch(() => {});
    };
    request.signal.addEventListener("abort", cleanup);

    return new Response(readable, {
      headers: {
        "Content-Type": "text/event-stream",
        "Cache-Control": "no-cache",
        "Connection": "keep-alive",
      },
    });
  }

  private async handleBroadcast(request: Request): Promise<Response> {
    const event = await request.json<SSEEvent>();
    const message = this.encoder.encode(`data: ${JSON.stringify(event)}\n\n`);

    const dead: WritableStreamDefaultWriter<Uint8Array>[] = [];
    for (const writer of this.connections) {
      try {
        await writer.write(message);
      } catch {
        dead.push(writer);
      }
    }
    dead.forEach((w) => this.connections.delete(w));

    return new Response("ok");
  }

  private async handleClose(): Promise<Response> {
    const event: SSEEvent = { type: "survey_closed" };
    const message = this.encoder.encode(`data: ${JSON.stringify(event)}\n\n`);

    for (const writer of this.connections) {
      try {
        await writer.write(message);
        await writer.close();
      } catch {}
    }
    this.connections.clear();
    return new Response("ok");
  }
}

有几点值得一提:

串行执行。Durable Objects 在一个实例内一次只处理一个请求。一次广播不会和另一次抢跑。不需要锁。

用中止信号来处理断开。清理不是靠心跳或错误检测循环,而是由 request.signal 驱动的。客户端一断开,中止事件就触发,写入器立刻从集合里移除。广播循环也会把写入时抛错的写入器剔掉,作为一道保险。

/close 端点。主讲人结束问卷时,管理路由会调用 /close,它向所有客户端发一个 survey_closed 事件,并干净地关掉这些连接。这让客户端代码可以做出反应——显示一句“问卷已结束”——而不是无声无息地丢掉连接。

首个事件。连接建立时,DO 立刻写一个 { type: "connected" } 事件。它向客户端确认这条 SSE 流是活的,也便于把“连接成功”和“连接挂住了”区分开。

路由到正确的房间

每份问卷的 Durable Object ID 都由它的问卷 ID 推导而来。在提交处理器(/api/respond)里:

const response = await createResponse(env.DB, surveyId, answers);

const doId = env.SURVEY_ROOM.idFromName(surveyId);
const stub = env.SURVEY_ROOM.get(doId);
await stub.fetch("http://do/broadcast", {
  method: "POST",
  body: JSON.stringify({ type: "new_response", response }),
  headers: { "Content-Type": "application/json" },
}).catch(() => {
  // 非致命:没有管理端在看时,这次广播就是空操作
});

idFromName() 是确定性的——同一个字符串在 Cloudflare 网络的任何地方,都会映射到同一个 DO 实例。这就是为什么你不需要任何协调层,也能拿到单实例保证。

广播是发完就不管的(.catch(() => {}))。如果没有管理面板连着,DO 就没有打开的写入器,这次广播也无害。

在流路由(/api/admin/[token]/stream)里:

const survey = await getSurveyByAdminToken(env.DB, token);
if (!survey) {
  return NextResponse.json({ error: "unauthorized" }, { status: 401 });
}

const doId = env.SURVEY_ROOM.idFromName(survey.id);
const stub = env.SURVEY_ROOM.get(doId);

return stub.fetch(new Request("http://do/connect", {
  signal: request.signal,
}));

这个路由先做认证,然后把请求的中止信号透传给 DO,好让它在浏览器离开页面时清理连接。

数据库:Cloudflare D1

问卷定义和答卷都存在 D1 里,也就是 Cloudflare 的边缘 SQLite。D1 在这儿很合适,因为问卷和答卷是关系型的,查询汇总结果用 SQL 是对的模型,而且这套模式简单到 SQLite 的那些限制根本不要紧。

对 SSE 流来说,D1 只读副本的最终一致性并不要紧——我们是通过 DO 直接把原始答卷数据广播出去,而不是每次提交后再去查一遍 D1。首屏加载时会查 D1 拿历史答卷,那没问题。

有一条实打实的约束:D1 会把写入串行化。每次提交都是一次串行的读加写,并发一高,这里就成了瓶颈。压测中,失败大约从 400 到 600 个虚拟用户开始出现。解法是用 Cloudflare Queues 和批量插入,把 D1 的写入从热路径上解耦出去——但那是以后的重构了。

部署:别扭的那部分

Cloudflare 的 next-on-pages 适配器会把 Next.js 应用转换成可以部署到 Pages 的形式,但你在应用里定义的 Durable Objects 需要作为一个独立的 Worker 部署。部署流程跑四步:

# 1. 构建 Next.js 并转换成 Cloudflare 格式
npx next build && npx @cloudflare/next-on-pages

# 2. 把 SurveyRoom DO 类注入编译好的 worker
node scripts/patch-worker.mjs

# 3. 部署 DO 伴随 Worker
wrangler deploy --config wrangler.worker.toml

# 4. 部署 Pages 项目
wrangler pages deploy .vercel/output/static --project-name instant-survey

patch-worker.mjs 这个脚本用 esbuild 编译 SurveyRoom.ts,并把产物拼在编译好的 _worker.js 包前面,好让 Pages worker 能导出这个 DO 类。它能用,但很脆——如果 Cloudflare 把 Pages 和 DO 的集成做好了,这是我第一个要重构的地方。

哪些做法奏效了,哪些没有

奏效的:

  • 对这个问题来说,DO 模型确实是对的抽象。一旦你把单实例路由这个保证内化了,实现就变得清楚而顺理成章。
  • SSE 跑在 Durable Objects 上非常稳。基于中止信号的清理,比心跳轮询更简单也更可靠。
  • D1 加 DO 合起来,把持久化和实时这一整套都覆盖了,不需要任何外部服务。整个应用都是 Cloudflare 原生的。

换做现在我会改的:

  • patch-worker 脚本是个权宜之计。把 DO 重构成一个有自己域名的完全独立 Worker、再从 Pages 应用代理过去,会更干净。
  • D1 的写入吞吐才是真正的扩展天花板。用 Queues 加批量插入把提交和 D1 写入解耦,能把上限推高不少。
  • 我应该更早在客户端加上显式的 Last-Event-ID 重连逻辑。SSE 连接是会断的,浏览器会自动重连,但有了序号,补回漏掉的事件就容易多了。

试一试

rifts.to 已经上线,而且免费。建一份问卷,拿上二维码,在你下一次团队会议或演讲里试试。两边都不需要注册。

相关工具

免费试用 rifts.to →
rifts.to
我如何用 Cloudflare Durable Objects 构建实时观众问卷 | rifts.to