•17 min read

Cloudflare Workers Durable Objects: Điều phối phân tán, WebSockets & Trạng thái biên

Cloudflare Workers Durable Objects: Điều phối phân tán, WebSockets & Trạng thái biên

Cloudflare Durable Objects (DOs) cung cấp khả năng lưu trữ và tính toán nhất quán mạnh mẽ, độ trễ thấp tại biên, cho phép xây dựng các ứng dụng phân tán toàn cầu, có trạng thái. Tài liệu này trình bày chi tiết các mẫu kiến trúc để tận dụng DOs nhằm quản lý trạng thái thời gian thực, tạo điều kiện phối hợp phân tán và xử lý hiệu quả các kết nối WebSocket có lưu lượng lớn.

Audio Briefing
0:00 / 0:00

Các Khái Niệm Cơ Bản về Durable Objects

Durable Object là một lớp thể hiện đơn lẻ tồn tại trên một trung tâm dữ liệu duy nhất của Cloudflare. Nó được xác định duy nhất bởi một DurableObjectId và có thể duy trì trạng thái giữa các yêu cầu. Tất cả các yêu cầu đến một thể hiện DO cụ thể đều được tuần tự hóa, đảm bảo tính nhất quán mạnh mẽ mà không cần cơ chế khóa rõ ràng. Việc tuần tự hóa này là một nguyên thủy cơ bản cho sự phối hợp phân tán.

Các Khái Niệm Cốt Lõi

  • Đảm bảo Thể hiện Đơn lẻ: Đối với một DurableObjectId nhất định, chỉ có một thể hiện của lớp DO tồn tại trên toàn cầu tại bất kỳ thời điểm nào.
  • Tính Nhất quán Mạnh mẽ: Tất cả các hoạt động trong một thể hiện DO đều là nguyên tử và được tuần tự hóa.
  • Lưu trữ Giao dịch: DOs cung cấp một kho khóa-giá trị không đồng bộ (state.storage) với ngữ nghĩa giao dịch. Các hoạt động trong một khối state.storage.transaction() duy nhất là nguyên tử.
  • Alarms: Một cơ chế để lên lịch một callback một lần hoặc định kỳ trong một DO, cho phép xử lý nền.
  • Ngủ đông WebSocket: DOs có thể quản lý hàng nghìn kết nối WebSocket không hoạt động với mức tiêu thụ tài nguyên tối thiểu.
Advertisement

Kiến trúc: Dịch vụ Chat Phân tán Thời gian thực

Chúng ta sẽ xây dựng một dịch vụ chat phân tán thời gian thực đơn giản để minh họa các khái niệm này. Mỗi phòng chat sẽ là một Durable Object. Người dùng kết nối đến một phòng thông qua WebSockets.

Cấu trúc Dự án

.
├── 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

Cấu hình wrangler.toml

Định nghĩa ràng buộc 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

Các kiểu dùng chung cho tin nhắn.

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

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

src/chatRoom.ts: Triển khai Durable Object

DO này quản lý trạng thái của phòng chat, xử lý các kết nối WebSocket và lưu trữ các tin nhắn gần đây.

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: Điểm vào Worker

Worker hoạt động như một bộ định tuyến, điều phối các yêu cầu đến thể hiện Durable Object thích hợp.

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 };

Các Mẫu & Tính năng Kiến trúc Chính

1. Ngủ đông WebSocket

Durable Objects vượt trội trong việc quản lý số lượng lớn các kết nối WebSocket không hoạt động. Khi một WebSocket không hoạt động (không có tin nhắn được gửi hoặc nhận), Cloudflare tự động "ngủ đông" nó. Thể hiện DO vẫn hoạt động, nhưng chi phí CPU liên quan đến kết nối thực tế bằng không. Chỉ khi một tin nhắn được gửi hoặc nhận, kết nối mới "thức dậy", phát sinh việc sử dụng CPU. Điều này cho phép một thể hiện DO duy nhất quản lý hàng chục nghìn kết nối đồng thời, không hoạt động.

Triển khai trong ChatRoom: Mảng sessions: WebSocket[] chứa các tham chiếu đến các kết nối WebSocket đang hoạt động. Runtime của Cloudflare xử lý việc ngủ đông một cách minh bạch. Các callback addEventListener chỉ được gọi khi có hoạt động.

2. Lưu trữ Giao dịch (Ngữ nghĩa giống SQLite)

API state.storage cung cấp một kho khóa-giá trị với tính nhất quán mạnh mẽ. Phương thức transaction() cho phép cập nhật nguyên tử, tất cả hoặc không có gì. Điều này rất quan trọng để duy trì tính toàn vẹn của dữ liệu.

Triển khai trong ChatRoom:

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

Điều này đảm bảo rằng this.messages được cập nhật hoàn toàn trong bộ nhớ hoặc không được cập nhật chút nào, ngăn chặn việc ghi một phần nếu DO gặp sự cố hoặc xảy ra lỗi trong quá trình giao dịch. Mặc dù không phải là một cơ sở dữ liệu quan hệ đầy đủ, transaction() cung cấp một nguyên thủy mạnh mẽ để quản lý trạng thái nhất quán.

3. Alarms cho Xử lý Nền

Alarms cho phép các tác vụ được lên lịch, một lần hoặc định kỳ trong một Durable Object. Điều này lý tưởng cho các công việc giống cron, dọn dẹp dữ liệu hoặc tổng hợp định kỳ.

Triển khai trong 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
}

Phương thức alarm() được runtime của Cloudflare gọi khi đến thời gian đã lên lịch. Nó được đảm bảo chạy tối đa một lần. Nếu DO bị loại bỏ, alarm vẫn tồn tại và sẽ kích hoạt khi DO được kích hoạt lại.

4. Phối hợp Phân tán

Đảm bảo thể hiện đơn lẻ của Durable Object là tính năng mạnh mẽ nhất của nó để phối hợp phân tán. Bất kỳ logic nào yêu cầu khóa toàn cục, một nguồn sự thật duy nhất hoặc truy cập tuần tự vào một tài nguyên đều có thể được đóng gói trong một DO.

Ví dụ: Cơ chế bầu chọn lãnh đạo. Một LeaderElectionDO có thể quản lý worker nào là lãnh đạo hiện tại, và các worker khác sẽ truy vấn DO này để xác định lãnh đạo. Trạng thái nội bộ của DO sẽ là nguồn có thẩm quyền.

Ví dụ: Giới hạn tốc độ. Một RateLimiterDO có thể theo dõi số lượng cuộc gọi API cho một người dùng hoặc IP nhất định, đảm bảo các giới hạn toàn cầu được thực thi trên tất cả các vị trí biên.

Đánh đổi & Cân nhắc

Tính năngDurable Objects (DO)Backend Truyền thống (ví dụ: EC2 + Redis)Serverless Functions (ví dụ: Lambda)
Tính nhất quánMạnh mẽ (thể hiện đơn lẻ, yêu cầu tuần tự)Cuối cùng (Redis), Mạnh mẽ (DB với khóa)Cuối cùng (nếu trạng thái chia sẻ là bên ngoài)
Độ trễCực thấp (tính toán biên)Thay đổi (tùy thuộc vào khu vực, số bước nhảy mạng)Thay đổi (khởi động lạnh, số bước nhảy mạng)
Quản lý trạng tháiKho KV giao dịch tích hợp, trong bộ nhớBên ngoài (Redis, DB), trong bộ nhớ (mỗi thể hiện)Bên ngoài (DB, S3, DynamoDB)
Khả năng mở rộngTự động theo từng đối tượng, mở rộng ngang của đối tượngNhóm mở rộng thủ công/tự động, sao chép trạng thái phức tạpTự động theo từng yêu cầu, không trạng thái theo thiết kế
Hỗ trợ WebSocketHạng nhất, ngủ đông cho các kết nối không hoạt độngYêu cầu máy chủ chuyên dụng, cân bằng tải phức tạpHạn chế, thường yêu cầu API Gateway + dịch vụ bên ngoài
Mô hình chi phíDựa trên yêu cầu, lưu trữ, thời gian CPU (bao gồm ngủ đông)Giờ thể hiện, truyền dữ liệu, phí dịch vụ được quản lýDựa trên yêu cầu, thời lượng, bộ nhớ
Độ phức tạpĐơn giản hơn cho logic biên có trạng thái, ít quản lý hạ tầngCao (hạ tầng, vận hành, mở rộng, HA)Trung bình (tích hợp, quản lý trạng thái)
Trường hợp sử dụngỨng dụng thời gian thực, trò chơi, chat, khóa phân tán, CRDTsỨng dụng web truyền thống, cơ sở dữ liệu phức tạpHướng sự kiện, xử lý hàng loạt, API
Advertisement

Những Vấn đề & Khắc phục sự cố trong Sản xuất

  1. DO Eviction và Khởi động Lạnh: Mặc dù DOs là "bền vững", VM cơ bản có thể bị loại bỏ và khởi tạo lại. Điều này có nghĩa là hàm tạo của DO của bạn sẽ chạy lại. Đảm bảo logic state.blockConcurrencyWhile() của bạn là bất biến và tải lại trạng thái cần thiết từ state.storage một cách hiệu quả. Nếu DO của bạn giữ trạng thái trong bộ nhớ lớn không được lưu trữ, nó sẽ bị mất khi bị loại bỏ.

    • Khắc phục: Lưu trữ trạng thái quan trọng vào state.storage. Sử dụng state.blockConcurrencyWhile để tải trạng thái đồng bộ trong quá trình xây dựng.
    • Triệu chứng: Mất dữ liệu không liên tục hoặc phản hồi ban đầu chậm sau thời gian không hoạt động.
  2. Lỗi state.storage.transaction(): Các giao dịch có thể thất bại do tranh chấp hoặc lỗi nội bộ. Callback async (txn) => { ... } sẽ được runtime của DO thử lại.

    • Khắc phục: Đảm bảo logic giao dịch của bạn là bất biến và xử lý việc thử lại một cách duyên dáng. Tránh các tác dụng phụ trong callback giao dịch mà không nên thực hiện lại.
    • Triệu chứng: Lỗi TransactionAborted trong nhật ký, hoặc trạng thái không nhất quán nếu không được xử lý đúng cách.
  3. Ngắt kết nối WebSocket: Khách hàng sẽ ngắt kết nối. DO của bạn phải xử lý những điều này một cách duyên dáng bằng cách xóa phiên khỏi danh sách hoạt động của nó. Các ngắt kết nối không được xử lý có thể dẫn đến rò rỉ bộ nhớ (giữ các đối tượng WebSocket đã đóng) hoặc cố gắng gửi tin nhắn đến các kết nối chết.

    • Khắc phục: Triển khai server.addEventListener('close', ...) và server.addEventListener('error', ...) để dọn dẹp sessions.
    • Triệu chứng: Lỗi Failed to send message to session, tăng mức sử dụng bộ nhớ theo thời gian, hoặc tin nhắn không đến được tất cả các khách hàng đang hoạt động.
  4. Lên lịch Alarm: Nếu một alarm được đặt cho một thời điểm trong quá khứ, nó sẽ kích hoạt ngay lập tức. Nếu bạn đặt một alarm và sau đó ngay lập tức đặt một alarm khác, cái thứ hai có thể ghi đè lên cái thứ nhất, hoặc cả hai có thể kích hoạt tùy thuộc vào thời gian chính xác và hành vi của runtime.

    • Khắc phục: Luôn kiểm tra await this.state.storage.getAlarm() trước khi đặt một alarm mới để tránh các alarm dư thừa hoặc xung đột, đặc biệt đối với các tác vụ định kỳ. Đảm bảo logic alarm của bạn mạnh mẽ để chống lại các kích hoạt lại tiềm năng.
    • Triệu chứng: Alarms kích hoạt bất ngờ hoặc không kích hoạt chút nào.
  5. Sử dụng CPU và Chi phí: Mặc dù ngủ đông WebSocket hiệu quả cho các kết nối không hoạt động, các kết nối đang hoạt động và tính toán nặng trong một DO sẽ tiêu thụ CPU. Một thể hiện DO duy nhất có tài nguyên hữu hạn. Nếu một DO duy nhất trở thành điểm nóng cho quá nhiều lưu lượng truy cập hoạt động hoặc tính toán phức tạp, nó có thể trở thành nút thắt cổ chai.

    • Khắc phục: Thiết kế ứng dụng của bạn để phân chia trạng thái trên nhiều DO (ví dụ: một DO cho mỗi phòng chat, một DO cho mỗi người dùng). Lập hồ sơ và tối ưu hóa các hoạt động tốn kém. Cân nhắc chuyển các tính toán nặng sang các Worker khác hoặc dịch vụ bên ngoài nếu nó không yêu cầu tính nhất quán mạnh mẽ trong DO.
    • Triệu chứng: Mức sử dụng CPU cao được báo cáo trong bảng điều khiển Cloudflare, tăng độ trễ cho các yêu cầu đến DO đó, hoặc lỗi DurableObjectStorageError: Too many concurrent operations.

Câu hỏi Thường gặp

  1. Durable Objects có thể giao tiếp với nhau không? Có. Một Durable Object có thể lấy một stub cho một Durable Object khác (ngay cả của một lớp khác) bằng cách sử dụng env.OTHER_DO_NAMESPACE.idFromName(id).get(id) và sau đó fetch() nó. Điều này cho phép các mẫu giao tiếp và phối hợp phức tạp giữa các DO.

  2. Làm cách nào để xử lý việc di chuyển lược đồ cho state.storage? state.storage là một kho khóa-giá trị, vì vậy việc di chuyển lược đồ là thủ công. Khi bạn thay đổi cấu trúc dữ liệu được lưu trữ dưới một khóa, bạn sẽ cần triển khai logic trong hàm tạo của DO hoặc phương thức fetch để phát hiện lược đồ cũ và di chuyển nó sang lược đồ mới. Điều này thường liên quan đến việc đọc dữ liệu cũ, chuyển đổi nó, ghi dữ liệu mới và tùy chọn xóa khóa cũ.

  3. Giới hạn lưu trữ Durable Object là gì? Mỗi Durable Object có thể lưu trữ tới 128KB dữ liệu trong state.storage. Giới hạn này áp dụng cho tổng kích thước của tất cả các cặp khóa-giá trị. Đối với dữ liệu lớn hơn, hãy cân nhắc lưu trữ các tham chiếu đến các đối tượng R2 hoặc bộ nhớ ngoài khác, và sử dụng DO cho siêu dữ liệu và phối hợp.

  4. Làm cách nào để đảm bảo Durable Object của tôi luôn khả dụng? Durable Objects vốn có tính khả dụng cao. Runtime của Cloudflare tự động di chuyển DOs giữa các trung tâm dữ liệu và khởi động lại chúng nếu máy cơ bản gặp sự cố. Bạn không cần triển khai logic chuyển đổi dự phòng rõ ràng. state.blockConcurrencyWhile() trong hàm tạo rất quan trọng để đảm bảo DO sẵn sàng phục vụ các yêu cầu sau bất kỳ kích hoạt hoặc di chuyển nào.

  5. Tôi có thể sử dụng Durable Objects cho các tính toán chạy dài không? Mặc dù DOs có thể chạy các tác vụ nền thông qua alarms, chúng không được thiết kế cho các tính toán chạy cực kỳ dài, tốn nhiều CPU. Mục đích chính là quản lý trạng thái và phối hợp. Nếu một tính toán mất quá nhiều thời gian, nó có thể vượt quá giới hạn thực thi của Cloudflare hoặc phát sinh chi phí CPU cao. Đối với xử lý hàng loạt nặng, hãy cân nhắc chuyển sang một dịch vụ tính toán chuyên dụng.

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 và Điện toán biên
serverless

Serverless và Điện toán biên

So sánh kiến trúc toàn diện giữa điện toán serverless và các runtime biên: hồ sơ độ trễ, giảm thiểu khởi động lạnh, trọng lực dữ liệu và đánh đổi chi phí.

Read more