TungDaDev's Blog

Transactional Outbox & Idempotent Consumer: Chấm dứt thảm họa Dual-Write trong Microservices

Photo 1526374965328 7f61d4dc18c5?auto=format&fit=crop&w=1600&q=80
Published on
/6 mins read/

Không có thứ gọi là Distributed Two-Phase Commit (2PC) vừa nhẹ nhàng vừa mở rộng tốt trên môi trường phân tán. Cách duy nhất để đạt được tính toàn vẹn dữ liệu giữa Database và Message Queue là dựa vào Transactional Outbox và Idempotent Processing.

1. Cơn ác mộng "Dual-Write" trong Microservices

Hãy xem đoạn mã Java kinh điển mà rất nhiều lập trình viên mới vào nghề hay viết:

@Transactional
public void placeOrder(OrderRequest request) {
    // 1. Lưu Order vào Database
    Order order = orderRepository.save(new Order(request));
 
    // 2. Bắn sự kiện lên Kafka để Payment Service xử lý
    kafkaTemplate.send("order-created-topic", new OrderCreatedEvent(order.getId()));
}

Trông rất gọn gàng và vô hại, nhưng đoạn code này chứa đựng 2 lỗi kiến trúc chí mạng:

Kịch bản thảm họa 1: Message bay đi nhưng DB Rollback

Lệnh kafkaTemplate.send() chạy thành công, message đã nằm trên broker Kafka. Nhưng ngay sau đó, một lỗi validation xảy ra hoặc database bị mất kết nối, @Transactional kích hoạt rollback. \rightarrow Hậu quả: Database không hề có đơn hàng nào, nhưng Payment Service đã nhận được message và trừ tiền trong tài khoản của khách hàng!

Kịch bản thảm họa 2: DB Commit thành công nhưng Kafka chết

Database đã commit thành công đơn hàng. Nhưng đúng lúc gửi sang Kafka thì mạng chập chờn hoặc Kafka cluster bị timeout. \rightarrow Hậu quả: Tiền khách hàng đã mất hoặc đơn đã tạo, nhưng không có email xác nhận, kho không xuất hàng, đơn hàng "chết lâm sàng".

Đây chính là bài toán Dual-Write Problem: Chúng ta không thể thực hiện một atomic commit trên hai hệ thống lưu trữ phân tán độc lập (RDBMS và Kafka) nếu không có sự hỗ trợ của các mẫu hình kiến trúc đặc thù.


2. Giải pháp: Transactional Outbox Pattern

Nguyên lý rất đơn giản nhưng hiệu quả tuyệt đối: Tận dụng chính Transaction ACID của cơ sở dữ liệu quan hệ (RDBMS).

Thay vì gửi trực tiếp sang Kafka:

  1. Mỗi khi ghi dữ liệu nghiệp vụ vào bảng orders, chúng ta ghi kèm một bản ghi sự kiện vào bảng outbox trong cùng một transaction duy nhất.
  2. Một tiến trình riêng biệt (Message Relay) sẽ đọc từ bảng outbox và đẩy sang Kafka.

2.1. Thiết kế bảng Outbox chuẩn mực

CREATE TABLE outbox_events (
    id UUID PRIMARY KEY,
    aggregate_type VARCHAR(255) NOT NULL, -- Ví dụ: 'ORDER'
    aggregate_id VARCHAR(255) NOT NULL,   -- order_id
    event_type VARCHAR(255) NOT NULL,     -- 'ORDER_CREATED'
    payload JSONB NOT NULL,               -- Dữ liệu chi tiết dạng JSON
    created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

2.2. Viết code nghiệp vụ an toàn

@Service
public class OrderService {
 
    private final OrderRepository orderRepository;
    private final OutboxEventRepository outboxRepository;
 
    @Transactional
    public Order placeOrder(OrderRequest request) {
        // 1. Lưu đơn hàng
        Order order = orderRepository.save(new Order(request));
 
        // 2. Lưu sự kiện vào outbox table (CÙNG TRANSACTION)
        OutboxEvent event = OutboxEvent.builder()
            .id(UUID.randomUUID())
            .aggregateType("ORDER")
            .aggregateId(order.getId().toString())
            .eventType("ORDER_CREATED")
            .payload(toJson(order))
            .build();
 
        outboxRepository.save(event);
        return order;
    }
}

Nếu bước 1 hoặc bước 2 lỗi, cả hai cùng rollback. Nếu thành công, cả hai chắc chắn được ghi vào đĩa!


3. Chuyển tiếp sự kiện từ Outbox lên Kafka: Polling vs Debezium CDC

Làm sao để đưa dữ liệu từ bảng outbox lên Kafka hiệu quả nhất?

Tiêu chíCách 1: Polling Publisher (@Scheduled)Cách 2: Transaction Log Tailing (Debezium CDC)
Độ phức tạp hạ tầngRất thấp (Chỉ cần vài chục dòng code Spring)Cần cụm Kafka Connect + Debezium
Tác động tới DatabaseGây tải I/O định kỳ (Query liên tục)Gần như bằng 0 (Đọc trực tiếp log file nhị phân)
Độ trễ (Latency)Phụ thuộc chu kỳ poll (vài giây)Gần như Real-time (< 50ms)
Khuyên dùngDự án vừa & nhỏ, tải thấpHệ thống lớn, tải hàng ngàn write/sec

4. Phía Consumer: Thiết kế Idempotent Consumer tuyệt đối

Vì cơ chế phân tán của Kafka và Debezium hoạt động theo chuẩn At-Least-Once Delivery (ít nhất một lần), message có thể bị gửi trùng lặp do:

  • Network disconnect khi commit offset Kafka.
  • Consumer xử lý xong nhưng crash trước khi gửi ACK về Kafka broker.

Do đó, Consumer bắt buộc phải có tính Idempotent (xử lý trùng không gây sai dữ liệu).

Triển khai Idempotent Consumer bằng PostgreSQL ON CONFLICT DO NOTHING:

CREATE TABLE processed_events (
    event_id UUID PRIMARY KEY,
    processed_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);
@Component
public class PaymentOrderConsumer {
 
    private final ProcessedEventRepository dedupRepo;
    private final PaymentService paymentService;
 
    @KafkaListener(topics = "order-created-topic")
    @Transactional
    public void handleOrderCreated(ConsumerRecord<String, String> record) {
        OrderCreatedEvent event = parseEvent(record.value());
 
        // 1. Kiểm tra và ghi nhận ID (Deduplication)
        boolean isNew = dedupRepo.insertIfNotExists(event.getEventId());
        if (!isNew) {
            log.warn("Duplicate event detected, skipping: {}", event.getEventId());
            return;
        }
 
        // 2. Thực thi nghiệp vụ an toàn
        paymentService.deductBalance(event.getCustomerId(), event.getAmount());
    }
}

5. Tổng kết Architecture Blueprint

Sự kết hợp giữa Transactional Outbox (ở phía Producer) và Idempotent Consumer (ở phía Consumer) tạo thành một bộ giáp hoàn chỉnh cho kiến trúc Event-Driven:

  1. Dữ liệu không bao giờ bị mất (No Data Loss).
  2. Không có hiện tượng dữ liệu ma do rollback (No Phantom Messages).
  3. Chống chịu hoàn hảo trước mọi lỗi mạng và trùng lặp bản tin (At-Least-Once biến thành Exactly-Once về mặt ngữ nghĩa).

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