+1

TỪ DATABASE ĐẾN KAFKA: CHIẾN LƯỢC LƯU TRỮ VÀ PHÂN PHỐI DỮ LIỆU

khi bước vào thế giới Microservices và kiến trúc hướng sự kiện (Event-Driven Architecture), việc lưu trữ dữ liệu xuống Database và đồng thời bắn thông báo lên Kafka tưởng chừng chỉ là hai dòng code đơn giản, nhưng thực chất lại là một trong những bài toán hóc búa nhất của hệ thống phân tán: Bài toán Ghi kép (Dual-Write Problem).

Dưới đây là bài viết mổ xẻ chi tiết cách thiết kế cấu trúc dữ liệu lưu trữ và chiến lược đẩy (publish) an toàn tuyệt đối lên Apache Kafka.

1. Sự Khác Biệt Giữa Cấu Trúc Lưu Trữ (State) và Cấu Trúc Tin Nhắn (Event)

Trước khi đẩy dữ liệu, bạn phải phân định rạch ròi hai khái niệm:

  • Data Storage (Database): Lưu trữ trạng thái hiện tại (Current State). Cấu trúc thường được chuẩn hóa (Normalized) thành các bảng quan hệ (MySQL/PostgreSQL) để tối ưu hóa truy vấn CRUD, giảm dư thừa dữ liệu. Dữ liệu ở đây có thể bị ghi đè (mutated).

  • Kafka Message (Event): Lưu trữ sự kiện đã xảy ra trong quá khứ (Historical Event). Cấu trúc tin nhắn thường phải phi chuẩn hóa (Denormalized) và gom thành một khối hoàn chỉnh để các service khác đọc là hiểu ngay mà không cần truy vấn ngược lại database. Tin nhắn trên Kafka là bất biến (immutable).

Ví dụ: Database của bạn lưu order_id = 1, status = 'paid'. Nhưng Event đẩy lên Kafka phải mang tên hành động: OrderPaidEvent, chứa toàn bộ snapshot dữ liệu lúc đó (tổng tiền, user_id, mã giao dịch).

2. "Tử Huyệt" Của Việc Gọi Kafka Trực Tiếp (Dual-Write Problem)

Cách tồi tệ nhất mà các lập trình viên mới thường làm là viết logic như sau trong Controller hoặc Service:

PHP

// ❌ CÁCH LÀM GÂY THẢM HỌA: Dual-Write
DB::beginTransaction();
$order = Order::create($data); // 1. Lưu vào Database
DB::commit();

// 2. Bắn lên Kafka
Kafka::publish('order_events', $order->toJson()); 

Chuyện gì sẽ xảy ra? Nếu Database commit thành công, nhưng tiến trình gọi lên Kafka bị lỗi mạng (Timeout) hoặc Kafka Broker đang down, hệ thống của bạn bị bất đồng bộ dữ liệu (Inconsistent). Database báo đã có đơn hàng, nhưng các service kho bãi, giao hàng (lắng nghe qua Kafka) không bao giờ nhận được thông báo!

3. Giải Pháp Tối Thượng: Transactional Outbox Pattern + CDC

Để giải quyết triệt để vấn đề trên, kiến trúc chuẩn mực (mà bạn có thể áp dụng với Debezium hoặc Poller) là Outbox Pattern.

Nguyên lý hoạt động: Thay vì bắn trực tiếp lên Kafka, chúng ta tạo thêm một bảng tên là outbox_events (chứa các sự kiện chuẩn bị gửi đi). Bạn gom lệnh tạo dữ liệu chính và lệnh tạo event vào cùng một Database Transaction.

SQL

-- Cấu trúc bảng Outbox
CREATE TABLE outbox_events (
    id UUID PRIMARY KEY,
    aggregate_type VARCHAR(255), -- Ví dụ: 'Order'
    aggregate_id VARCHAR(255),   -- Ví dụ: '105'
    event_type VARCHAR(255),     -- Ví dụ: 'OrderCreated'
    payload JSONB,               -- Cấu trúc dữ liệu đẩy lên Kafka
    created_at TIMESTAMP
);

Luồng thực thi trong Code:

PHP

// ✅ CÁCH LÀM CHUẨN: Transactional Outbox
DB::transaction(function () use ($data) {
    // 1. Lưu nghiệp vụ chính
    $order = Order::create($data); 
    
    // 2. Lưu Event vào bảng Outbox CÙNG MỘT GIAO DỊCH
    OutboxEvent::create([
        'aggregate_type' => 'Order',
        'aggregate_id' => $order->id,
        'event_type' => 'OrderCreated',
        'payload' => $order->toJson()
    ]);
});
// Đảm bảo 100% nếu DB lưu thành công thì Outbox cũng lưu thành công.

Làm sao để dữ liệu từ Outbox lên được Kafka?

  • Cách 1 (Change Data Capture - CDC): Sử dụng các công cụ mạnh mẽ như Debezium (chạy bằng Kafka Connect) để đọc trực tiếp file Binlog/WAL của MySQL/PostgreSQL. Cứ có dòng mới thêm vào bảng outbox_events, Debezium sẽ tự động bắt lấy và ném lên Kafka với độ trễ tính bằng mili-giây.

  • Cách 2 (Polling Worker): Viết một background worker (Cronjob hoặc Daemon) liên tục quét bảng outbox_events mỗi giây, lấy dữ liệu gửi lên Kafka, gửi thành công thì xóa dòng đó đi hoặc đánh dấu status = 'published'.

4. Chuẩn Hóa Cấu Trúc Message (Payload Structure)

Khi đẩy dữ liệu (payload) từ Outbox lên Kafka, việc chọn định dạng rất quan trọng:

  • JSON: Dễ đọc, dễ debug nhưng kích thước nặng và không ràng buộc kiểu dữ liệu chặt chẽ (Schema validation).

  • Protobuf (Protocol Buffers) / Avro: Định dạng nhị phân siêu nhẹ, tốc độ phân giải nhanh gấp nhiều lần JSON. Đặc biệt, nó kết hợp với Schema Registry, buộc các microservices (dù viết bằng Go, PHP hay Node.js) phải tuân thủ đúng một hợp đồng dữ liệu (contract), chống hoàn toàn lỗi một bên đổi tên trường dữ liệu làm bên kia sập hệ thống.

💡 Lời Kết

Việc cấu trúc lưu trữ và publish lên Kafka không chỉ là bài toán chuyển đổi định dạng, mà là bài toán đảm bảo tính toàn vẹn (Consistency). Bằng cách áp dụng Transactional Outbox kết hợp với cơ chế CDC như Debezium và định dạng Protobuf, hệ thống phân tán của bạn sẽ sở hữu một luồng chảy dữ liệu bền bỉ, không bao giờ rơi rớt tin nhắn dù có gặp sự cố mạng hay sập server giữa chừng.


All rights reserved

Viblo
Hãy đăng ký một tài khoản Viblo để nhận được nhiều bài viết thú vị hơn.
Đăng kí