Durable Objects — Workers上のステートフル&リアルタイム
Durable Objectsを活用したステートフル開発:WebSocket、ゲームサーバー、分散協調、組み込みSQLiteストレージ。
Durable Objects(DO)は、Cloudflare Workersにステートフルな協調レイヤーを提供するソリューションです。各DOはグローバルに唯一のインスタンスとして動作し、強整合性(Strong Consistency)を保ちながら状態を保持します。
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/>1つのインスタンスへルーティング"| 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ごとに1つの調停インスタンス)"]
subgraph DO1["Chat Room DO (ID: room-123)"]
Mem["メモリ内ステート & 接続中WebSocket"]
SQL[("組み込み SQLite / KV ストレージ")]
Alarm["アラーム / タイマー"]
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 |
|---|---|
| ステートレス | ステートフル |
| リクエストごとに複数インスタンス | 1つのIDに対して世界で1つのインスタンス |
| メモリ状態を保持しない | メモリ(RAM)+ストレージに状態を保持 |
| 無限に自動スケール | シャーディング(複数の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 — 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 });
}
}
}typescriptstorage.sql による SQL クエリ#
Durable Object 内で直接 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 });
}
}typescriptアラーム — DO のスケジュールタスク#
sequenceDiagram
autonumber
actor User as ユーザー / クライアント
participant DO as Reminder DO
participant Storage as ctx.storage [アラームキュー]
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 & State からプレイヤーを削除<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 の移行管理#
DO はクラスやスキーマの移行をダウンタイムなしでサポートします:
// wrangler.jsonc
{
"durable_objects": {
"bindings": [
{
"name": "COUNTER",
"class_name": "Counter",
"migration": "new_tag"
}
]
}
}typescriptマイグレーションタイプ:
new_tag— 新規クラスタグnew_classes— 複数クラスの追加rename— クラス名変更transfer— ネームスペース間のデータ移行
パターン — シャーディング(Sharding)#
単一の Durable Object インスタンスは物理的に1箇所で動作します。水平スケールを行うにはキーでシャーディングします:
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 に委譲します。キーによるシャーディングで無限のスケールアウトが可能です。