TungDaDev's Blog

Figma: scale PostgreSQL với Horizontal Sharding

Figma postgresql sharding.webp
Published on
/12 mins read/

Khi một cơ sở dữ liệu quan hệ nguyên khối chạm trần phần cứng của instance đám mây lớn nhất, lựa chọn giữ lại PostgreSQL và triển khai Horizontal Sharding đòi hỏi phải giải quyết triệt để ba bài toán kỹ thuật: ranh giới colocation của dữ liệu, kiểm soát phân luồng truy vấn ở tầng proxy, và ngăn chặn bùng nổ kết nối mạng.

# giới hạn vật lý của cụm PostgreSQL nguyên khối

Trước khi phân mảnh (sharding), toàn bộ trạng thái tài liệu, quyền truy cập và dữ liệu cộng tác thời gian thực của Figma tập trung tại một cụm AWS Aurora PostgreSQL nguyên khối (monolithic). Dù đã phân tách tải đọc sang hàng loạt Read Replicas, nút Primary Writer vẫn phải gánh toàn bộ lưu lượng ghi (INSERT, UPDATE, DELETE) của hàng chục triệu người dùng hoạt động đồng thời.

Figma đã nâng cấp cụm Primary lên instance lớn nhất mà AWS cung cấp ở thời điểm đó (128 vCPU, 1024 GB RAM, trần 64,000 EBS IOPS). Tuy nhiên, hệ thống vẫn đối mặt với các giới hạn kiến trúc không thể scale dọc:

  • Bão hòa băng thông CPU và IOPS: Lưu lượng ghi dồn dập khiến CPU của Primary duy trì trên ngưỡng 80%, trong khi dung lượng ghi nhật ký Write-Ahead Logging (WAL) vượt quá 100 MB/s, tiến sát ngưỡng nghẽn I/O của hạ tầng lưu trữ phân tán Aurora.
  • Hiện tượng Replication Lag kéo dài: Khối lượng WAL khổng lồ từ một writer duy nhất làm các replica xử lý redo log không kịp, dẫn đến dữ liệu đọc trên replica bị trễ hàng giây. Điều này phá vỡ tính nhất quán của các tính năng cộng tác trực tiếp.
  • Áp lực bộ nhớ do tiến trình backend: PostgreSQL sử dụng mô hình đa tiến trình (process-per-connection). Mỗi kết nối chiếm khoảng 5 MB đến 10 MB bộ nhớ chia sẻ. Khi hàng nghìn container backend đồng thời mở kết nối, bộ nhớ dành cho shared_buffers và caching trang bảng (page cache) bị thu hẹp nghiêm trọng.

Figma cân nhắc giữa hai hướng đi: Chuyển dịch toàn bộ hệ thống sang NoSQL phân tán (DynamoDB, Spanner, Cassandra) hoặc xây dựng kiến trúc Horizontal Sharding trên chính PostgreSQL. Họ chọn phương án thứ hai để bảo toàn mô hình quan hệ, ngữ nghĩa giao dịch ACID trên từng tài liệu, và hệ sinh thái công cụ đã kiểm chứng qua nhiều năm vận hành.


# thiết kế Colocated Tables và phân cấp Shard Key

Nguyên tắc cốt lõi của Sharding trong RDBMS là: Tuyệt đối tránh Cross-Shard JOINs và Distributed 2PC (Two-Phase Commit). Nếu một câu lệnh truy vấn buộc phải lấy dữ liệu từ hai máy chủ vật lý khác nhau rồi kết hợp lại qua mạng, độ trễ hệ thống sẽ tăng từ vài mili-giây lên hàng trăm mili-giây.

Để giải quyết triệt để điều này, Figma áp dụng mô hình Colocation: Toàn bộ các thực thể dữ liệu có quan hệ mật thiết với nhau phải cùng nằm trên một node vật lý (Shard).

# chuẩn hóa Shard Key cấp gốc

Mọi thực thể nghiệp vụ cốt lõi trong Figma đều xoay quanh hai ranh giới phân vùng tự nhiên:

  • Tài liệu (File): Đối với dữ liệu bản vẽ, layer, comment, lịch sử phiên bản (canvas_nodes, comments, versions), Shard Key là file_id.
  • Tổ chức (Organization): Đối với dữ liệu phòng ban, dự án, quyền thành viên (teams, projects, permissions), Shard Key là org_id.

Mọi bảng con nằm trong ranh giới file bắt buộc phải chứa cột file_id, ngay cả khi bảng đó có khóa chính riêng. Ví dụ: bảng comments không chỉ có comment_id, mà bắt buộc phải mang theo file_id để router biết chính xác vị trí lưu trữ.

# xử lý khóa ngoại và định danh toàn cục

  • Khóa ngoại (Foreign Keys): RDBMS không thể tự động thực thi ràng buộc khóa ngoại giữa hai cơ sở dữ liệu độc lập. Do đó, Figma loại bỏ Foreign Keys giữa các bảng nằm ngoài cùng ranh giới Shard Key, chuyển trách nhiệm kiểm tra tính toàn vẹn tham chiếu về tầng ứng dụng (Application Level Validation).
  • Sinh ID phân tán: Khóa chính tự tăng (BIGSERIAL) trên từng node sẽ gây xung đột ID trên toàn cụm. Hệ thống chuyển sang sử dụng bộ sinh định danh 64-bit toàn cục (dựa trên thuật toán tương tự Twitter Snowflake/KSUID), kết hợp timestamp, ID của node worker, và chuỗi tuần tự để đảm bảo tính duy nhất và có thứ tự theo thời gian.

# phân định Global Tables và Sharded Tables

Không phải bảng nào cũng có thể chia theo file_id hay org_id. Figma phân loại bảng thành hai nhóm rõ rệt:

  1. Sharded Tables: Chiếm hơn 95% dung lượng và tải ghi, được băm đều vào N Shards vật lý bằng hàm băm nhất quán: Shard Index = MurmurHash3(file_id) % N
  2. Global Tables: Các bảng nhỏ, ít ghi nhưng cần truy cập toàn cục (như bảng xác thực người dùng users, danh mục tính năng toàn cầu feature_flags) được lưu trữ tại một cụm Global Shard riêng biệt.

# kiến trúc Query Router: Phân tích cú pháp SQL AST

Để các kỹ sư phần mềm tại Figma tiếp tục sử dụng ORM và viết truy vấn SQL tiêu chuẩn mà không cần tự tính toán vị trí shard trong từng dòng code, Figma phát triển một lớp trung gian (Database Proxy / Query Router).

Lớp Proxy này đóng vai trò như một bộ biên dịch và điều phối truy vấn:

  1. Phân tích cây cú pháp (SQL AST Parsing): Khi nhận một câu truy vấn, proxy bóc tách cú pháp SQL thành cây trừu tượng (Abstract Syntax Tree). Nó xác định danh sách các bảng tham gia và quét điều kiện WHERE để tìm giá trị của Shard Key (file_id hoặc org_id).
  2. Định tuyến đơn shard (Single-Shard Routing): Sau khi trích xuất được giá trị file_id, proxy tra cứu bảng phân vùng (Shard Map) và chuyển tiếp toàn bộ truy vấn nguyên bản đến chính xác shard đích. Quá trình này diễn ra với độ trễ cộng thêm dưới 0.5 mili-giây.
  3. Thực thi ranh giới an toàn (Query Guardrails): Nếu một kỹ sư vô tình viết câu lệnh SELECT * FROM comments WHERE author_id = ? (thiếu file_id), proxy sẽ lập tức chặn lại và ném lỗi DisallowedScatterQueryException. Điều này ngăn chặn việc broadcast câu lệnh tới toàn bộ 50+ shards vật lý cùng lúc, bảo vệ cụm database khỏi hiện tượng nghẽn quạt (fan-out explosion).

# giải quyết bài toán bùng nổ kết nối M x N

Khi chia cơ sở dữ liệu thành N shards vật lý và hệ thống backend có M máy chủ ứng dụng (pods), số lượng kết nối tối đa có thể phát sinh là:

Tổng kết nối = M × N × Pool Size

Nếu có 2,000 backend pods kết nối tới 48 shards với kích thước pool là 10, tổng số kết nối mở vào database là 2,000 \times 48 \times 10 = 960,000 kết nối. Con số này vượt xa ngưỡng chịu đựng của PostgreSQL kernel.

Figma giải quyết vấn đề bằng kiến trúc gom nhóm hai tầng:

  • Giữa ứng dụng và Proxy: Giao tiếp qua các kết nối nhẹ (lightweight connection pool) hoặc gRPC streaming.
  • Giữa Proxy và PostgreSQL: Triển khai Transaction-level Pooling (tương tự cơ chế của PgBouncer). Proxy chỉ mượn một kết nối thực sự tới Shard PostgreSQL khi câu lệnh SQL bắt đầu thực thi transaction và trả lại kết nối vào pool ngay khi lệnh COMMIT hoặc ROLLBACK kết thúc.
  • Nhờ cơ chế này, số lượng kết nối thực tế trên mỗi Shard PostgreSQL được khống chế ở mức ổn định từ 200 đến 400 kết nối, tối ưu hóa triệt để bộ nhớ đệm và giảm chi phí chuyển đổi ngữ cảnh (context switching) của CPU.

# quy trình cắt tải dữ liệu live (Zero-Downtime Migration)

Chuyển đổi hàng chục Terabyte dữ liệu từ một cơ sở dữ liệu đang hoạt động 24/7 sang hàng chục shards mới mà không được làm gián đoạn trải nghiệm của hàng triệu người dùng trực tuyến đòi hỏi một quy trình 4 bước chặt chẽ:

# 1. nạp dữ liệu nền bằng Change Data Capture (CDC)

Hệ thống sử dụng cơ chế PostgreSQL Logical Replication thông qua plugin pgoutput để đọc luồng thay đổi nhị phân từ WAL của Monolithic DB. Các worker xử lý sẽ giải mã bản ghi, trích xuất Shard Key, tính toán shard đích và nạp vào cụm Shards mới mà không gây khóa bảng trên cụm chính.

# 2. ghi kép và đọc ngầm đối chiếu (Shadow Verification)

Trước khi chuyển đổi lưu lượng thực, tầng ứng dụng kích hoạt cơ chế ghi kép bất đồng bộ (Asynchronous Dual-Writing) và đọc ngầm (Shadow Reading):

  • Mỗi truy vấn đọc từ người dùng sẽ được gửi đồng thời tới cả Monolithic DB và Shard mới.
  • Kết quả từ Monolithic DB được trả về cho người dùng, trong khi một tiến trình ngầm so sánh từng byte dữ liệu giữa hai nguồn.
  • Mọi trường hợp lệch dữ liệu (data discrepancy) đều kích hoạt cảnh báo để đội ngũ kỹ thuật hiệu chỉnh thuật toán đồng bộ.

# 3. cắt tải trong tích tắc (Atomic Cutover)

Khi độ trễ đồng bộ giữa Monolithic và Shards đạt mức xấp xỉ 0:

  • Hệ thống tạm thời chuyển Monolithic DB sang chế độ READ_ONLY trong khoảng 1-2 giây.
  • Worker CDC xả hết những bản ghi WAL cuối cùng vào cụm Shards.
  • Proxy cập nhật cấu hình: chuyển toàn bộ lưu lượng đọc và ghi sang cụm Sharded Cluster.
  • Hệ thống hoạt động bình thường trên kiến trúc mới mà không làm mất bất kỳ giao dịch nào của người dùng.

# tổng kết các nguyên tắc kiến trúc

Hành trình scale 100 lần của Figma để lại những kinh nghiệm kỹ thuật giá trị cho các hệ thống tải lớn:

Khía cạnhCách tiếp cận truyền thốngKiến trúc Sharding của Figma
Lựa chọn lưu trữVội vã chuyển sang NoSQL phân tánGiữ nguyên PostgreSQL, phân mảnh theo chiều ngang
Phân vùng dữ liệuPhân mảnh ngẫu nhiên theo ID dòngColocation Design: Nhóm toàn bộ dữ liệu phụ thuộc quanh một Root Key
Ràng buộc toàn vẹnPhụ thuộc vào Foreign Key trong DBChuyển kiểm tra tính toàn vẹn lên tầng ứng dụng
Kiểm soát truy vấnĐể ứng dụng tự kết nối tự doSử dụng SQL AST Proxy có guardrails chặn truy vấn đa shard
Quản lý kết nốiMở kết nối trực tiếp từ BackendDùng Transaction Pooling khống chế số lượng kết nối thực

Việc chia nhỏ cơ sở dữ liệu quan hệ theo chiều ngang đòi hỏi đầu tư lớn vào hạ tầng proxy và công cụ di chuyển dữ liệu, nhưng mang lại khả năng mở rộng gần như tuyến tính mà vẫn bảo toàn tính toàn vẹn của mô hình quan hệ phức tạp.


Chỉ là những ghi chép cá nhân với hy vọng mang lại chút giá trị. Nếu thấy hữu ích, đừng ngại chia sẻ cho bạn bè & đồng nghiệp nhé!

Happy coding 😎 👍🏻 🚀 🔥.