Durable Objects — Workers 上的有状态实时应用
使用 Durable Objects 构建有状态应用:WebSocket、游戏服务器、协同调度以及内置 SQLite 存储。
Durable Objects(DO)是 Cloudflare Workers 的有状态协同解决方案。每个 DO 拥有全球唯一的运行实例 —— 状态持久存在,并具备强一致性保证。
flowchart LR
subgraph EdgeWorkers["无状态 Workers(全球边缘节点)"]
direction TB
W_Tokyo["Worker 东京"]
W_London["Worker 伦敦"]
W_US["Worker 美国"]
end
subgraph DO_Instance["Durable Object(全球唯一实例)"]
direction TB
ID["ID: 'user-room-42'"]
subgraph Memory["RAM(内存状态)"]
State["活跃连接与状态数据"]
end
subgraph Storage["持久化存储"]
SQL[("SQLite / KV Storage")]
end
ID --> Memory
Memory <--> Storage
end
W_Tokyo -->|"所有相同 ID 的请求<br/>路由至单一实例"| DO_Instance
W_London -->|"强一致性<br/>(杜绝竞态条件)"| DO_Instance
W_US -->|"状态持续驻留"| DO_Instance
Workers vs Durable Objects#
flowchart TB
subgraph Clients["客户端与浏览器"]
C1["用户 A (WebSocket)"]
C2["用户 B (WebSocket)"]
C3["用户 C (HTTP REST)"]
end
subgraph Edge["Cloudflare 全球边缘(无状态 Workers)"]
W1["Worker(边缘节点 1)"]
W2["Worker(边缘节点 2)"]
end
subgraph DO_Cluster["Durable Objects(每个 ID 对应单一协调实例)"]
subgraph DO1["Chat Room DO (ID: room-123)"]
Mem["内存状态 & 活跃 WebSocket"]
SQL[("内置 SQLite / KV 存储")]
Alarm["定时器 / Alarms"]
Mem <--> SQL
Mem <--> Alarm
end
subgraph DO2["Chat Room DO (ID: room-456)"]
Mem2["内存状态"]
SQL2[("SQLite 存储")]
end
end
C1 <-->|"边缘连接"| W1
C2 <-->|"边缘连接"| W2
C3 -->|"HTTP 请求"| W1
W1 <-->|"路由至 ID: room-123"| DO1
W2 <-->|"路由至 ID: room-123"| DO1
W1 -.->|"路由至 ID: room-456"| DO2
| Workers | Durable Objects |
|---|---|
| 无状态(Stateless) | 有状态(Stateful) |
| 多个请求产生多个实例 | 每个 ID 对应全球唯一实例 |
| 内存不保留状态 | 在内存 + 存储中持久保留状态 |
| 无限自动水平扩展 | 通过分片(Sharding 多 DO)扩展 |
| 分布在 330+ 个数据中心 | 在单一位置执行(支持热迁移) |
| --- | --- |
聊天室 — 基础示例#
sequenceDiagram
autonumber
actor ClientA as 用户 Alice
actor ClientB as 用户 Bob
participant Worker as 无状态 Worker 路由
participant DO as ChatRoom [Durable Object]
Note over ClientA, DO: 初始化 WebSocket 连接
ClientA->>Worker: GET /?room=general&name=Alice (Upgrade: WebSocket)
Worker->>DO: stub.fetch(request) [idFromName("general")]
DO->>DO: server.accept(), 保存会话 "Alice"
DO-->>ClientA: 101 Switching Protocols (WebSocket 已连接)
DO--)ClientB: 广播 "Alice joined the chat"
Note over ClientA, DO: 实时消息广播
ClientA->>DO: WS 消息: "Hello room!"
DO->>DO: broadcast("Alice: Hello room!", sender=Alice)
DO--)ClientB: WS 消息: "Alice: Hello room!"
import { DurableObject } from 'cloudflare:workers';
interface Env {
CHAT_ROOM: DurableObjectNamespace;
}
export class ChatRoom extends DurableObject {
private sessions = new Map<string, WebSocket>();
async fetch(request: Request): Promise<Response> {
const url = new URL(request.url);
const name = url.searchParams.get('name') || 'anonymous';
// WebSocket upgrade
const pair = new WebSocketPair();
const [client, server] = Object.values(pair);
server.accept();
this.sessions.set(name, server);
// Broadcast welcome
this.broadcast(`${name} joined the chat`);
server.addEventListener('message', (event) => {
this.broadcast(`${name}: ${event.data}`, server);
});
server.addEventListener('close', () => {
this.sessions.delete(name);
this.broadcast(`${name} left`);
});
return new Response(null, { status: 101, webSocket: client });
}
private broadcast(message: string, sender?: WebSocket) {
const data = JSON.stringify({ message, timestamp: Date.now() });
this.sessions.forEach((ws, name) => {
if (ws !== sender && ws.readyState === WebSocket.OPEN) {
ws.send(data);
}
});
}
}
// Worker — router 到 DO
export default {
async fetch(request: Request, env: Env): Promise<Response> {
const url = new URL(request.url);
const roomId = url.searchParams.get('room') || 'default';
const id = env.CHAT_ROOM.idFromName(roomId);
const stub = env.CHAT_ROOM.get(id);
return stub.fetch(request);
},
};typescriptSQLite Storage — 内置存储引擎#
每个 Durable Object 均拥有独立的 SQLite 存储:
flowchart TD
Req["请求: GET /increment /get /reset"] --> Match{url.pathname}
Match -->|"/increment"| Inc["storage.get('count')<br/>storage.put('count', count + 1)"]
Match -->|"/get"| Get["storage.get('count')"]
Match -->|"/reset"| Res["storage.put('count', 0)"]
Match -->|其他| Err["404 Not Found"]
Inc --> Resp["返回 JSON count"]
Get --> Resp
Res --> Resp
export class Counter extends DurableObject {
private storage: DurableObjectStorage;
constructor(ctx: DurableObjectState, env: Env) {
super(ctx, env);
this.storage = ctx.storage;
}
async fetch(request: Request): Promise<Response> {
const url = new URL(request.url);
switch (url.pathname) {
case '/increment':
const count = (await this.storage.get<number>('count')) || 0;
await this.storage.put('count', count + 1);
return Response.json({ count: count + 1 });
case '/get':
const current = (await this.storage.get<number>('count')) || 0;
return Response.json({ count: current });
case '/reset':
await this.storage.put('count', 0);
return Response.json({ count: 0 });
default:
return new Response('Not found', { status: 404 });
}
}
}typescript使用 storage.sql 执行 SQL 查询#
直接在 Durable Objects 中执行原生 SQL:
flowchart TD
Req["HTTP 请求"] --> Init["表不存在则创建: CREATE TABLE IF NOT EXISTS todos (...)"]
Init --> Router{HTTP 方法}
Router -->|"GET"| Q1["SELECT * FROM todos ORDER BY id DESC"]
Router -->|"POST"| Q2["INSERT INTO todos (title) VALUES (?)"]
Router -->|"PUT"| Q3["UPDATE todos SET completed = ? WHERE id = ?"]
Router -->|"DELETE"| Q4["DELETE FROM todos WHERE id = ?"]
Router -->|其他| Q5["405 Method Not Allowed"]
Q1 --> Out1["JSON: todos 列表"]
Q2 --> Out2["201 Created: 新创建的 todo"]
Q3 --> Out3["JSON: { success: true }"]
Q4 --> Out4["JSON: { success: true }"]
export class TodoApp extends DurableObject {
async fetch(request: Request): Promise<Response> {
const sql = this.ctx.storage.sql;
// Create table
sql.exec(`
CREATE TABLE IF NOT EXISTS todos (
id INTEGER PRIMARY KEY AUTOINCREMENT,
title TEXT NOT NULL,
completed INTEGER DEFAULT 0
)
`);
const url = new URL(request.url);
// List
if (request.method === 'GET') {
const result = sql.exec('SELECT * FROM todos ORDER BY id DESC');
return Response.json(result.toArray());
}
// Create
if (request.method === 'POST') {
const { title } = await request.json();
const result = sql.exec('INSERT INTO todos (title) VALUES (?)', title);
return Response.json({ id: result.lastRowId, title, completed: false }, { status: 201 });
}
// Update
if (request.method === 'PUT') {
const { id, completed } = await request.json();
sql.exec('UPDATE todos SET completed = ? WHERE id = ?', completed, id);
return Response.json({ success: true });
}
// Delete
if (request.method === 'DELETE') {
const id = url.searchParams.get('id');
sql.exec('DELETE FROM todos WHERE id = ?', id);
return Response.json({ success: true });
}
return new Response('Method not allowed', { status: 405 });
}
}typescriptAlarms — DO 定时任务与调度#
sequenceDiagram
autonumber
actor User as 用户 / 客户端
participant DO as Reminder DO
participant Storage as ctx.storage [Alarm 队列]
participant Webhook as 外部 Webhook
Note over DO, Storage: 1. 注册定时器 / 闹钟
User->>DO: GET /?delay=5000
DO->>Storage: setAlarm(futureTimestamp)
DO-->>User: {"message": "Alarm set for 5000ms"}
Note over Storage, Webhook: 2. 超时唤醒执行(唤醒 DO)
Storage-->>DO: 触发 alarm() 处理器
DO->>Webhook: POST https://hooks.example.com/notify { event: 'alarm_fired' }
DO->>Storage: setAlarm(nextHourTimestamp) [设置下次循环定时]
export class Reminder extends DurableObject {
constructor(ctx: DurableObjectState, env: Env) {
super(ctx, env);
ctx.blockConcurrencyWhile(async () => {
const alarm = await ctx.storage.getAlarm();
if (alarm) {
console.log('Recovered alarm:', new Date(alarm).toISOString());
}
});
}
async fetch(request: Request): Promise<Response> {
const url = new URL(request.url);
const delayMs = parseInt(url.searchParams.get('delay') || '5000');
// Set alarm sau delayMs
await this.ctx.storage.setAlarm(Date.now() + delayMs);
return Response.json({ message: `Alarm set for ${delayMs}ms` });
}
async alarm() {
// Được gọi khi alarm hết giờ
console.log('⏰ Alarm fired!');
// Gửi notification qua WebSocket nếu có
// Hoặc call webhook
await fetch('https://hooks.example.com/notify', {
method: 'POST',
body: JSON.stringify({ event: 'alarm_fired', time: Date.now() }),
});
// Set alarm lại nếu cần
await this.ctx.storage.setAlarm(Date.now() + 3600000);
}
}typescript多人在线游戏服务器#
flowchart TD
WS["WebSocket 消息事件"] --> EventType{msg.type}
EventType -->|"join"| Join["1. 将玩家会话加入 Map<br/>2. 初始化状态: x=0, y=0<br/>3. 广播: type: 'players'"]
EventType -->|"move"| Move["1. 更新坐标: player.x += dx, player.y += dy<br/>2. 广播: type: 'move'"]
EventType -->|"shoot"| Shoot["1. 计算子弹方向与角度<br/>2. 广播: type: 'shoot'"]
Close["WebSocket 断开事件"] --> Leave["1. 从 Map 和状态中移除玩家<br/>2. 广播: type: 'leave'"]
Join --> BroadcastAll["向房间内所有连接玩家广播更新后的状态"]
Move --> BroadcastAll
Shoot --> BroadcastAll
Leave --> BroadcastAll
interface Player {
id: string;
x: number;
y: number;
score: number;
}
export class GameRoom extends DurableObject {
private players = new Map<string, WebSocket>();
private state: Player[] = [];
async fetch(request: Request): Promise<Response> {
const pair = new WebSocketPair();
const [client, server] = Object.values(pair);
server.accept();
server.addEventListener('message', (event) => {
const msg = JSON.parse(event.data as string);
switch (msg.type) {
case 'join':
this.players.set(msg.playerId, server);
this.state.push({ id: msg.playerId, x: 0, y: 0, score: 0 });
this.broadcast({ type: 'players', data: this.state });
break;
case 'move':
const player = this.state.find(p => p.id === msg.playerId);
if (player) {
player.x += msg.dx;
player.y += msg.dy;
this.broadcast({ type: 'move', playerId: msg.playerId, x: player.x, y: player.y });
}
break;
case 'shoot':
this.broadcast({
type: 'shoot',
playerId: msg.playerId,
x: msg.x,
y: msg.y,
angle: msg.angle,
});
break;
}
});
server.addEventListener('close', () => {
this.state = this.state.filter(p => !this.players.has(p.id));
this.players.delete(msg.playerId);
this.broadcast({ type: 'leave', playerId: msg.playerId });
});
return new Response(null, { status: 101, webSocket: client });
}
private broadcast(message: object) {
const data = JSON.stringify(message);
this.players.forEach((ws) => {
if (ws.readyState === WebSocket.OPEN) {
ws.send(data);
}
});
}
}typescript迁移管理 — 跨版本迁移 DO#
Durable Objects 支持零停机在不同类和架构间进行迁移:
// wrangler.jsonc
{
"durable_objects": {
"bindings": [
{
"name": "COUNTER",
"class_name": "Counter",
"migration": "new_tag"
}
]
}
}typescript迁移类型:
new_tag— 新类标签new_classes— 添加多个类rename— 类重命名transfer— 命名空间之间的数据迁移
设计模式 — 分片(Sharding)#
每个 Durable Object 实例在同一时刻仅运行于单个物理位置。为了实现横向扩展,通常通过 Key 进行分片:
flowchart LR
Req["请求: userId"] --> Formula["shardId = floor(userId / 1000)"]
Formula --> ShardMap{分片路由}
ShardMap -->|"userId: 0 - 999"| S0["DO 实例: shard-0"]
ShardMap -->|"userId: 1000 - 1999"| S1["DO 实例: shard-1"]
ShardMap -->|"userId: 2000 - 2999"| S2["DO 实例: shard-2"]
ShardMap -->|"userId: N - N+999"| Sn["DO 实例: shard-N"]
// Shard user data theo userId
function getUserStub(userId: number, env: Env): DurableObjectStub {
const shardId = Math.floor(userId / 1000); // 1000 user mỗi shard
const id = env.USER_STORE.idFromName(`shard-${shardId}`);
return env.USER_STORE.get(id);
}
// Handler
export default {
async fetch(request: Request, env: Env): Promise<Response> {
const url = new URL(request.url);
const userId = parseInt(url.searchParams.get('userId') || '0');
const stub = getUserStub(userId, env);
return stub.fetch(request);
},
};typescript总结#
Durable Objects 完美解决了无服务器架构中的状态管理难题:
- 实时通信:WebSocket 连接结合全球强一致性内存状态
- 游戏服务器:低延迟多人在线状态协同
- 分布式协同:全局分布式锁、选主、速率限制
- 内置 SQLite:无需维护外部数据库即可进行全功能关系型查询
架构设计建议:无状态请求由 Workers 处理,需要一致性状态协同的场景交由 Durable Objects。通过 Key 进行分片以实现无限水平扩展。