TungDaDev's Blog

Discord: từ Cassandra sang ScyllaDB

Discord scylladb.webp
Published on
/9 mins read/

Khi dữ liệu tăng trưởng từ hàng triệu lên hàng nghìn tỷ bản ghi, một quyết định kiến trúc cơ sở dữ liệu đúng đắn không chỉ quyết định sự sống còn của dịch vụ, mà còn tiết kiệm hàng triệu USD chi phí hạ tầng mỗi năm.

# tiến trình kiến trúc lưu trữ tại Discord

Discord ra mắt năm 2015 với một cụm MongoDB duy nhất. Khi số lượng người dùng bùng nổ, toàn bộ dữ liệu tin nhắn nhanh chóng vượt dung lượng RAM, đẩy disk I/O vào trạng thái nghẽn liên tục.

Năm 2017, Discord di chuyển sang Apache Cassandra — cơ sở dữ liệu NoSQL Wide-Column phân tán. Cassandra đã đồng hành cùng Discord qua nhiều cột mốc tăng trưởng. Tuy nhiên, khi dữ liệu chạm mốc hàng nghìn tỷ tin nhắn trên 177 node, những rào cản vật lý của máy ảo Java và cơ chế I/O truyền thống bắt đầu bộc lộ.


# mô hình dữ liệu tin nhắn và cấu trúc Partition

Trong Cassandra và ScyllaDB, cách thiết kế Partition Key quyết định trực tiếp việc dữ liệu có phân bổ đều trên cluster hay không. Discord mô hình hóa bảng tin nhắn theo cấu trúc:

CREATE TABLE messages (
    channel_id bigint,
    bucket int,
    message_id bigint,
    author_id bigint,
    content text,
    PRIMARY KEY ((channel_id, bucket), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);
  • Partition Key ((channel_id, bucket)): Xác định node nào trong cluster sẽ lưu trữ dữ liệu.
  • Clustering Key (message_id): Định đoạt thứ tự sắp xếp vật lý của các dòng tin nhắn bên trong từng partition. Do message_id tại Discord sử dụng Snowflake ID (chứa timestamp nhúng), các tin nhắn mới nhất luôn nằm ở đầu để tối ưu cho việc truy vấn đọc tin nhắn gần đây.

# tại sao cần kỹ thuật Bucket?

Nếu chỉ dùng channel_id làm partition key, một kênh chat lớn với hàng chục triệu tin nhắn sẽ phình to thành một partition khổng lồ (vài Gigabyte). Khi một partition vượt quá 100MB, hiệu năng của Cassandra tụt dốc thảm hại. Do đó, Discord chia nhỏ kênh chat thành các bucket theo thời gian hoặc theo khối 10 ngày để giữ kích thước mỗi partition luôn dưới ngưỡng an toàn.


# các điểm nghẽn của Apache Cassandra ở quy mô nghìn tỷ

Mặc dù được thiết kế để mở rộng theo chiều ngang, Cassandra bắt đầu bộc lộ 3 vấn đề cố hữu khi dữ liệu mở rộng đến hàng nghìn tỷ bản ghi:

# hiện tượng Stop-The-World (STW) của JVM Garbage Collection

Cassandra chạy trên JVM. Với lưu lượng hàng trăm nghìn request/giây, lượng object ngắn hạn sinh ra trên Heap là cực lớn. Dù đội ngũ Discord đã tinh chỉnh từ CMS, G1GC đến ZGC, các đợt Stop-The-World kéo dài từ 200ms đến vài giây vẫn xảy ra ngẫu nhiên khi GC thu gom rác trên các node chịu tải cao, kéo độ trễ P99/P99.9 tăng đột biến.

# vấn đề Tombstone Overwhelming

Khi người dùng xóa tin nhắn hoặc tin nhắn hết hạn, Cassandra không xóa dữ liệu trên đĩa ngay lập tức mà ghi một điểm đánh dấu gọi là Tombstone. Khi câu lệnh SELECT quét qua một khoảng thời gian chứa nhiều tin nhắn đã xóa, coordinator phải đọc hàng chục nghìn tombstone trước khi tìm thấy dữ liệu sống:

-- Nếu khoảng quét chứa quá 100,000 tombstones:
TombstoneOverwhelmingException: Nodes scanned over 100001 tombstones...

Hiện tượng này tạo ra áp lực rác khổng lồ lên JVM Heap và làm tê liệt thread đọc của Cassandra.

# Compaction Storm và tranh chấp I/O

Cơ chế Log-Structured Merge-tree (LSM-tree) đòi hỏi liên tục gộp các SSTable nhỏ thành SSTable lớn. Vào giờ cao điểm, khi lượng ghi dồn dập, quá trình Compaction tranh chấp băng thông ổ cứng với các tác vụ đọc ghi của người dùng, khiến hàng đợi I/O bị nghẽn (I/O queuing).


# ScyllaDB: Kiến trúc Thread-Per-Core và Seastar Engine

Để giải quyết triệt để các vấn đề trên, Discord chuyển đổi sang ScyllaDB — cơ sở dữ liệu tương thích hoàn toàn với giao thức CQL của Cassandra, nhưng được viết lại hoàn toàn bằng C++ trên nền tảng framework bất đồng bộ Seastar.

Những cải tiến kiến trúc mang tính quyết định của ScyllaDB:

  1. Mô hình Shared-Nothing (Thread-per-core):

    • Mỗi core CPU được gán cứng một luồng thực thi (pinned thread), sở hữu riêng một phần RAM và hàng đợi I/O.
    • Không có lock, không có mutex, không có tranh chấp bộ nhớ giữa các nhân CPU.
    • Giao tiếp giữa các core sử dụng hàng đợi vòng không khóa (Lock-free Ring Buffer) với chi phí CPU gần như bằng 0.
  2. Quản lý bộ nhớ trực tiếp không qua JVM:

    • Loại bỏ 100% hiện tượng Stop-The-World của Garbage Collection.
    • Bộ nhớ cache của ScyllaDB được quản lý ở tầng người dùng (Userspace), tối ưu hóa chính xác cho cấu trúc hàng dữ liệu NoSQL thay vì phụ thuộc vào OS Page Cache.
  3. Direct I/O và kiểm soát I/O Scheduler:

    • Giao tiếp trực tiếp với ổ cứng NVMe qua Linux AIO / io_uring, bỏ qua lớp đệm của nhân hệ điều hành.
    • Tự động điều tiết băng thông: Nếu lưu lượng đọc của người dùng tăng cao, I/O Scheduler sẽ tự động giảm tốc độ của tiến trình Compaction để ưu tiên phục vụ request.

# chiến lược di chuyển dữ liệu live (Zero-Downtime Migration)

Di chuyển hàng nghìn tỷ bản ghi đang phục vụ trực tuyến đòi hỏi một kế hoạch kiểm soát rủi ro chặt chẽ. Discord đã triển khai dịch vụ di chuyển dữ liệu viết bằng Rust theo quy trình 4 giai đoạn:

  1. Dual Writing: Tất cả tin nhắn mới được ghi đồng thời vào cả Cassandra và ScyllaDB. Lỗi ghi sang ScyllaDB được ghi log để bù đắp mà không làm gián đoạn request của người dùng.
  2. Historical Data Migration: Worker Rust chia nhỏ không gian hash 64-bit token của cluster thành các dải nhỏ (Token Ranges), quét tuần tự và chép dữ liệu lịch sử sang ScyllaDB với cơ chế điều tiết tốc độ tự động (Rate Limiting).
  3. Shadow Reading: Mỗi truy vấn đọc dữ liệu từ Cassandra được gửi kèm một bản sao bất đồng bộ sang ScyllaDB để kiểm tra tính nhất quán từng trường dữ liệu.
  4. Cutover: Khi tỷ lệ sai lệch giữa hai hệ thống đạt mức 0%, toàn bộ lưu lượng được chuyển hướng chính thức sang ScyllaDB.

# kết quả thực tế

Sau khi hoàn tất di chuyển, cụm database của Discord ghi nhận sự cải thiện rõ rệt:

Chỉ số đo đạcApache Cassandra cũScyllaDB mớiMức độ cải thiện
Quy mô cluster177 nodes72 nodes📉 Giảm 59% số node
Độ trễ đọc P99400ms – 1.2s (dao động mạnh)< 15ms (ổn định)⚡ Nhanh hơn ~40 lần
Độ trễ ghi P99100ms – 250ms< 5ms⚡ Nhanh hơn ~30 lần
Thời gian dừng do GCXảy ra liên tục mỗi ngày0ms (hoàn toàn triệt tiêu)🎯 Độ tin cậy tuyệt đối

TIP

Bài học quan trọng từ case study này không phải là "Cassandra kém cỏi", mà là việc hiểu rõ giới hạn phần cứng: Khi bài toán mở rộng đến ngưỡng hàng chục Terabyte trên mỗi node với hàng trăm nghìn QPS, kiến trúc Thread-per-core bằng C++ của ScyllaDB khai thác phần cứng hiện đại hiệu quả hơn đáng kể so với kiến trúc đa luồng chia sẻ bộ nhớ trên JVM.


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 😎 👍🏻 🚀 🔥.