•20 min read

Cloudflare Workers Durable Objects:分散協調、WebSockets、エッジステート

Cloudflare Workers Durable Objects:分散協調、WebSockets、エッジステート

Cloudflare Durable Objects(DO)は、エッジで強力な一貫性と低遅延のストレージおよびコンピューティングを提供し、グローバルに分散されたステートフルなアプリケーションの構築を可能にします。このドキュメントでは、DOを活用してリアルタイムの状態を管理し、分散協調を促進し、大量のWebSocket接続を効率的に処理するためのアーキテクチャパターンについて詳しく説明します。

Audio Briefing
0:00 / 0:00

Durable Objectsの基礎

Durable Objectは、単一のCloudflareデータセンター上に存在する単一インスタンスのクラスです。DurableObjectIdによって一意に識別され、リクエスト間で状態を維持できます。特定のDOインスタンスへのすべてのリクエストはシリアル化され、明示的なロックメカニズムなしで強力な一貫性を保証します。このシリアル化は、分散協調のための基本的なプリミティブです。

コアコンセプト

  • 単一インスタンス保証(Single-Instance Guarantee): 特定のDurableObjectIdに対して、DOクラスのインスタンスは常にグローバルに1つだけ存在します。
  • 強力な一貫性(Strong Consistency): DOインスタンス内のすべての操作はアトミックかつシリアル化されます。
  • トランザクションストレージ(Transactional Storage): DOは、トランザクションセマンティクスを持つ非同期キーバリューストア(state.storage)を提供します。単一のstate.storage.transaction()ブロック内の操作はアトミックです。
  • アラーム(Alarms): DO内でワンショットまたは定期的なコールバックをスケジュールするメカニズムで、バックグラウンド処理を可能にします。
  • WebSocketハイバネーション(WebSocket Hibernation): DOは、最小限のリソース消費で数千のアイドル状態のWebSocket接続を管理できます。
Advertisement

アーキテクチャ:リアルタイム分散チャットサービス

これらの概念を説明するために、簡略化されたリアルタイム分散チャットサービスを構築します。各チャットルームはDurable Objectになります。ユーザーはWebSocketを介してルームに接続します。

プロジェクト構造

.
├── src
│   ├── index.ts          # Worker entry point, routes requests to DOs
│   ├── chatRoom.ts       # Durable Object for a chat room
│   └── types.ts          # Shared types
├── wrangler.toml         # Cloudflare Workers configuration
└── package.json

wrangler.toml設定

Durable Objectバインディングを定義します。

name = "chat-service"
main = "src/index.ts"
compatibility_date = "2024-01-01"

[[durable_objects.bindings]]
name = "CHAT_ROOM"
class_name = "ChatRoom"
script_name = "chat-service" # Optional: if DO is in a different worker

src/types.ts

メッセージの共有型。

export interface ChatMessage {
  sender: string;
  message: string;
  timestamp: number;
}

export interface WebSocketMessage {
  type: 'chat' | 'join' | 'leave';
  payload: ChatMessage | { user: string };
}

src/chatRoom.ts:Durable Objectの実装

このDOは、チャットルームの状態を管理し、WebSocket接続を処理し、最近のメッセージを永続化します。

import { ChatMessage, WebSocketMessage } from './types';

// Maximum number of messages to store in memory and persist
const MAX_MESSAGES = 100;
// Key for storing messages in Durable Object storage
const MESSAGES_STORAGE_KEY = 'messages';

export class ChatRoom {
  private state: DurableObjectState;
  private env: Env; // Environment variables, if any
  private sessions: WebSocket[] = []; // Active WebSocket connections
  private messages: ChatMessage[] = []; // In-memory message buffer

  constructor(state: DurableObjectState, env: Env) {
    this.state = state;
    this.env = env;

    // Restore messages from storage on DO activation
    this.state.blockConcurrencyWhile(async () => {
      const storedMessages = await this.state.storage.get<ChatMessage[]>(MESSAGES_STORAGE_KEY);
      if (storedMessages) {
        this.messages = storedMessages;
      }
      console.log(`ChatRoom ${this.state.id.toString()} initialized with ${this.messages.length} messages.`);
    });
  }

  // Handles incoming HTTP requests, primarily WebSocket upgrades
  async fetch(request: Request): Promise<Response> {
    const url = new URL(request.url);

    if (url.pathname === '/websocket') {
      // Upgrade the request to a WebSocket connection
      const { 0: client, 1: server } = new WebSocketPair();

      this.handleWebSocket(server);

      return new Response(null, { status: 101, webSocket: client });
    } else if (url.pathname === '/messages') {
      // Return recent messages via HTTP for initial load or REST clients
      return new Response(JSON.stringify(this.messages), {
        headers: { 'Content-Type': 'application/json' },
      });
    }

    return new Response('Not Found', { status: 404 });
  }

  // Handles the lifecycle of a single WebSocket connection
  private async handleWebSocket(server: WebSocket) {
    server.accept();
    this.sessions.push(server);

    // Send historical messages to the new client
    server.send(JSON.stringify({ type: 'history', payload: this.messages }));

    server.addEventListener('message', async (event) => {
      try {
        const message: WebSocketMessage = JSON.parse(event.data as string);
        if (message.type === 'chat') {
          const chatMessage: ChatMessage = {
            ...message.payload as ChatMessage,
            timestamp: Date.now(),
          };
          this.broadcastMessage(chatMessage);
          await this.storeMessage(chatMessage);
        } else if (message.type === 'join') {
          const user = (message.payload as { user: string }).user;
          this.broadcastSystemMessage(`${user} joined the room.`);
        }
      } catch (err) {
        console.error('WebSocket message parsing error:', err);
        server.send(JSON.stringify({ type: 'error', payload: 'Invalid message format.' }));
      }
    });

    server.addEventListener('close', async (event) => {
      console.log(`WebSocket closed: ${event.code} ${event.reason}`);
      this.sessions = this.sessions.filter(s => s !== server);
      // Optionally broadcast leave message
    });

    server.addEventListener('error', (err) => {
      console.error('WebSocket error:', err);
      server.close(1011, 'Internal Error');
    });
  }

  // Broadcasts a message to all connected clients
  private broadcastMessage(message: ChatMessage) {
    const wsMessage: WebSocketMessage = { type: 'chat', payload: message };
    const serializedMessage = JSON.stringify(wsMessage);
    this.sessions.forEach(session => {
      try {
        session.send(serializedMessage);
      } catch (err) {
        console.error('Failed to send message to session:', err);
        // Consider removing broken sessions here
      }
    });
  }

  // Broadcasts a system message (e.g., join/leave)
  private broadcastSystemMessage(text: string) {
    const systemMessage: ChatMessage = { sender: 'System', message: text, timestamp: Date.now() };
    const wsMessage: WebSocketMessage = { type: 'chat', payload: systemMessage };
    const serializedMessage = JSON.stringify(wsMessage);
    this.sessions.forEach(session => {
      try {
        session.send(serializedMessage);
      } catch (err) {
        console.error('Failed to send system message to session:', err);
      }
    });
  }

  // Stores a message in memory and persists it to Durable Object storage
  private async storeMessage(message: ChatMessage) {
    this.messages.push(message);
    if (this.messages.length > MAX_MESSAGES) {
      this.messages.shift(); // Keep only the latest MAX_MESSAGES
    }

    // Persist messages transactionally
    await this.state.storage.transaction(async (txn) => {
      txn.put(MESSAGES_STORAGE_KEY, this.messages);
    });

    // Set an alarm to prune old messages or perform other background tasks
    // This demonstrates using alarms for cron-like functionality
    const currentAlarm = await this.state.storage.getAlarm();
    if (currentAlarm === null || currentAlarm < Date.now() + 60 * 1000) { // Set alarm for 1 minute from now
      await this.state.storage.setAlarm(Date.now() + 60 * 1000);
    }
  }

  // Alarm handler for background tasks
  async alarm() {
    console.log(`Alarm triggered for ChatRoom ${this.state.id.toString()}. Pruning messages.`);
    // Example: Prune messages older than 24 hours
    const twentyFourHoursAgo = Date.now() - (24 * 60 * 60 * 1000);
    const initialLength = this.messages.length;
    this.messages = this.messages.filter(msg => msg.timestamp > twentyFourHoursAgo);

    if (this.messages.length < initialLength) {
      // Only persist if messages were actually pruned
      await this.state.storage.transaction(async (txn) => {
        txn.put(MESSAGES_STORAGE_KEY, this.messages);
      });
      console.log(`Pruned ${initialLength - this.messages.length} messages.`);
    }

    // Optionally reschedule the alarm for the next pruning cycle
    await this.state.storage.setAlarm(Date.now() + 24 * 60 * 60 * 1000); // Reschedule for 24 hours later
  }
}

// Define the environment interface for type safety
interface Env {
  CHAT_ROOM: DurableObjectNamespace;
}

src/index.ts:Workerのエントリポイント

Workerはルーターとして機能し、適切なDurable Objectインスタンスにリクエストをディスパッチします。

import { ChatRoom } from './chatRoom';

interface Env {
  CHAT_ROOM: DurableObjectNamespace;
}

export default {
  async fetch(request: Request, env: Env, ctx: ExecutionContext): Promise<Response> {
    const url = new URL(request.url);
    const roomId = url.pathname.split('/')[1] || 'default-room'; // Extract room ID from path

    // Get a Durable Object ID for the room
    const id = env.CHAT_ROOM.idFromName(roomId);

    // Get the Durable Object stub
    const stub = env.CHAT_ROOM.get(id);

    // Forward the request to the Durable Object
    return stub.fetch(request);
  },
};

// Export the Durable Object class for Wrangler to discover
export { ChatRoom };

主要なアーキテクチャパターンと機能

1. WebSocketハイバネーション

Durable Objectsは、多数のアイドル状態のWebSocket接続を管理するのに優れています。WebSocketがアイドル状態(メッセージの送受信がない)の場合、Cloudflareは自動的に「ハイバネート」します。DOインスタンスはアクティブなままですが、接続に関連するCPUコストは実質的にゼロです。メッセージが送受信された場合にのみ、接続は「ウェイクアップ」し、CPU使用量が発生します。これにより、単一のDOインスタンスで数万の同時アイドル接続を管理できます。

ChatRoomでの実装: sessions: WebSocket[]配列は、アクティブなWebSocket接続への参照を保持します。Cloudflareのランタイムはハイバネーションを透過的に処理します。addEventListenerコールバックは、アクティビティが発生した場合にのみ呼び出されます。

2. トランザクションストレージ(SQLiteのようなセマンティクス)

state.storage APIは、強力な一貫性を持つキーバリューストアを提供します。transaction()メソッドは、アトミックなオールオアナッシングの更新を可能にします。これは、データの整合性を維持するために不可欠です。

ChatRoomでの実装:

await this.state.storage.transaction(async (txn) => {
  txn.put(MESSAGES_STORAGE_KEY, this.messages);
});

これにより、DOがクラッシュしたり、トランザクション中にエラーが発生したりした場合に、this.messagesがストレージ内で完全に更新されるか、まったく更新されないかのどちらかになり、部分的な書き込みを防ぎます。完全なリレーショナルデータベースではありませんが、transaction()は一貫性のある状態管理のための強力なプリミティブを提供します。

3. バックグラウンド処理のためのアラーム

アラームは、Durable Object内でスケジュールされたワンショットまたは定期的なタスクを可能にします。これは、cronのようなジョブ、データクリーンアップ、または定期的な集計に最適です。

ChatRoomでの実装:

// Setting an alarm
await this.state.storage.setAlarm(Date.now() + 60 * 1000); // 1 minute from now

// Alarm handler
async alarm() {
  // ... background task logic ...
  await this.state.storage.setAlarm(Date.now() + 24 * 60 * 60 * 1000); // Reschedule
}

alarm()メソッドは、スケジュールされた時刻が到来したときにCloudflareランタイムによって呼び出されます。最大1回実行されることが保証されています。DOがエビクトされた場合でも、アラームは永続化され、DOが再アクティブ化されたときにトリガーされます。

4. 分散協調

Durable Objectの単一インスタンス保証は、分散協調のための最も強力な機能です。グローバルロック、単一の信頼できる情報源、またはリソースへのシリアル化されたアクセスを必要とするロジックは、DO内にカプセル化できます。

例: リーダー選出メカニズム。LeaderElectionDOは、どのワーカーが現在のリーダーであるかを管理でき、他のワーカーはこのDOにクエリを実行してリーダーを決定します。DOの内部状態が権威ある情報源となります。

例: レート制限。RateLimiterDOは、特定のユーザーまたはIPのAPI呼び出し数を追跡し、すべてのエッジロケーションでグローバルな制限が適用されるようにできます。

トレードオフと考慮事項

機能Durable Objects (DO)従来のバックエンド (例: EC2 + Redis)サーバーレス関数 (例: Lambda)
一貫性強力 (単一インスタンス、シリアル化されたリクエスト)結果整合性 (Redis)、強力 (ロック付きDB)結果整合性 (共有状態が外部の場合)
レイテンシ非常に低い (エッジコンピューティング)可変 (リージョン、ネットワークホップに依存)可変 (コールドスタート、ネットワークホップ)
状態管理組み込みのトランザクションKVストア、インメモリ外部 (Redis, DB)、インメモリ (インスタンスごと)外部 (DB, S3, DynamoDB)
スケーリングオブジェクトごとの自動、オブジェクトの水平スケーリング手動/オートスケーリンググループ、複雑な状態レプリケーションリクエストごとの自動、設計上ステートレス
WebSocketサポートファーストクラス、アイドル接続のハイバネーション専用サーバー、複雑なロードバランシングが必要制限あり、API Gateway + 外部サービスが必要な場合が多い
コストモデルリクエストベース、ストレージ、CPU時間 (ハイバネーション含む)インスタンス時間、データ転送、マネージドサービス料金リクエストベース、実行時間、メモリ
複雑さステートフルなエッジロジックがよりシンプル、インフラ管理が少ない高 (インフラ、運用、スケーリング、HA)中程度 (統合、状態管理)
ユースケースリアルタイムアプリ、ゲーム、チャット、分散ロック、CRDTs従来のWebアプリ、複雑なデータベースイベント駆動型、バッチ処理、API
Advertisement

本番環境での注意点とトラブルシューティング

  1. DOのエビクションとコールドスタート: DOは「耐久性がある」とはいえ、基盤となるVMはエビクトされて再インスタンス化される可能性があります。これは、DOのコンストラクタが再度実行されることを意味します。state.blockConcurrencyWhile()ロジックが冪等であり、state.storageから必要な状態を効率的に再ロードするようにしてください。DOが永続化されていない大きなインメモリ状態を保持している場合、エビクション時に失われます。

    • 修正: 重要な状態をstate.storageに永続化します。state.blockConcurrencyWhileを使用して、構築中に状態を同期的にロードします。
    • 症状: 非アクティブ期間後の断続的なデータ損失または初期応答の遅延。
  2. state.storage.transaction()の失敗: 競合や内部エラーによりトランザクションが失敗する可能性があります。async (txn) => { ... }コールバックはDOランタイムによって再試行されます。

    • 修正: トランザクションロジックが冪等であり、再試行を適切に処理するようにしてください。再実行すべきではないトランザクションコールバック内で副作用を避けてください。
    • 症状: ログにTransactionAbortedエラーが表示されるか、正しく処理されない場合に状態が不整合になる。
  3. WebSocketの切断: クライアントは切断されます。DOは、アクティブなリストからセッションを削除することで、これらを適切に処理する必要があります。未処理の切断は、メモリリーク(閉じられたWebSocketオブジェクトを保持し続ける)や、切断された接続へのメッセージ送信の試行につながる可能性があります。

    • 修正: server.addEventListener('close', ...)とserver.addEventListener('error', ...)を実装して、sessionsをクリーンアップします。
    • 症状: Failed to send message to sessionエラー、時間の経過とともにメモリ使用量の増加、またはアクティブなすべてのクライアントにメッセージが届かない。
  4. アラームのスケジューリング: 過去の時刻にアラームが設定されている場合、すぐにトリガーされます。アラームを設定した直後に別のアラームを設定すると、2番目のアラームが最初のものを上書きしたり、正確なタイミングとランタイムの動作によっては両方がトリガーされたりする可能性があります。

    • 修正: 特に定期的なタスクの場合、冗長または競合するアラームを避けるために、新しいアラームを設定する前に常にawait this.state.storage.getAlarm()を確認してください。アラームロジックが潜在的な再トリガーに対して堅牢であることを確認してください。
    • 症状: アラームが予期せず発火したり、まったく発火しなかったりする。
  5. CPU使用量とコスト: WebSocketハイバネーションはアイドル接続には効率的ですが、アクティブな接続とDO内の重い計算はCPUを消費します。単一のDOインスタンスには有限のリソースがあります。単一のDOがアクティブなトラフィックや複雑な計算のホットスポットになりすぎると、ボトルネックになる可能性があります。

    • 修正: 複数のDOに状態をシャーディングするようにアプリケーションを設計します(例: チャットルームごとに1つのDO、ユーザーごとに1つのDO)。高価な操作をプロファイルして最適化します。DO内で強力な一貫性を必要としない場合は、重い計算を他のWorkerまたは外部サービスにオフロードすることを検討してください。
    • 症状: Cloudflareダッシュボードで報告される高いCPU使用量、そのDOへのリクエストのレイテンシの増加、またはDurableObjectStorageError: Too many concurrent operationsエラー。

よくある質問

  1. Durable Objectsは互いに通信できますか? はい。Durable Objectは、env.OTHER_DO_NAMESPACE.idFromName(id).get(id)を使用して別のDurable Object(異なるクラスのものでも)のスタブを取得し、それをfetch()できます。これにより、複雑なDO間の通信と協調パターンが可能になります。

  2. state.storageのスキーマ移行はどのように処理しますか? state.storageはキーバリューストアであるため、スキーマ移行は手動です。キーの下に保存されているデータの構造を変更する場合、DOのコンストラクタまたはfetchメソッドで、古いスキーマを検出し、新しいスキーマに移行するロジックを実装する必要があります。これには通常、古いデータを読み取り、変換し、新しいデータを書き込み、オプションで古いキーを削除する作業が含まれます。

  3. Durable Objectストレージの制限は何ですか? 各Durable Objectは、state.storageに最大128KBのデータを保存できます。この制限は、すべてのキーと値のペアの合計サイズに適用されます。より大きなデータの場合は、R2オブジェクトまたは他の外部ストレージへの参照を保存し、DOをメタデータと協調に使用することを検討してください。

  4. Durable Objectが常に利用可能であることをどのように保証しますか? Durable Objectsは本質的に高可用性です。Cloudflareのランタイムは、DOをデータセンター間で自動的に移行し、基盤となるマシンが故障した場合は再起動します。明示的なフェイルオーバーロジックを実装する必要はありません。コンストラクタのstate.blockConcurrencyWhile()は、DOがアクティベーションまたは移行後にリクエストを処理する準備ができていることを保証するために不可欠です。

  5. Durable Objectsを長時間実行される計算に使用できますか? DOはアラームを介してバックグラウンドタスクを実行できますが、非常に長時間実行されるCPU集中型の計算向けには設計されていません。主な目的は状態管理と協調です。計算に時間がかかりすぎると、Cloudflareの実行制限に達したり、高いCPUコストが発生したりする可能性があります。重いバッチ処理の場合は、専用のコンピューティングサービスにオフロードすることを検討してください。

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement
ServerlessとEdgeComputing
serverless

ServerlessとEdgeComputing

Serverlessコンピューティングとedge runtimeのアーキテクチャを、レイテンシプロファイル、コールドスタート緩和、データグラビティ、コストトレードオフの観点から包括的に比較します。

Read more