Mười ba service dùng chung một nền tảng chuỗi cà phê, và service nào cũng giữ một niềm tin riêng về đơn hàng. Checkout giữ chính đơn hàng. Loyalty nợ khách điểm tích luỹ. Inventory phải giữ hàng. Reporting cộng doanh thu. Khoảnh khắc bốn bên hết khớp nhau, retry bao nhiêu lần cũng vô ích — vì lúc đó không ai nói được bên nào mới đúng.
Đây là sự cố đẩy tôi đến transactional outbox, và nếu chỉ được giữ đúng một pattern, tôi giữ cái này.
Một lần ghi mà hoá ra là hai
Bản ngây thơ trông rất hợp lý:
await prisma.order.create({ data: order })
await kafka.send({ topic: 'order.placed', messages: [{ value: json }] })
Hai lần ghi. Hai hệ thống. Không transaction nào bao được cả hai.
Mọi thứ nằm giữa hai dòng đó đều là chỗ để hỏng. Pod bị evict. Broker đang rebalance. Request timeout trong khi broker thực ra đã nhận message. Postgres commit xong rồi process chết trước khi chạy tới dòng thứ hai.
Kết quả: đơn hàng thì có, sự kiện thì không. Loyalty không bao giờ cộng điểm. Inventory không bao giờ giữ hàng. Khách nhận xác nhận đặt hàng, rồi một tuần sau nhận email báo hết hàng.
Vì sao retry không cứu được
Phản xạ thường thấy là bọc lệnh publish trong retry. Nó không giải quyết được, và nên nói rõ vì sao.
Retry chỉ chạy được khi process còn sống. Mà những lỗi thực sự gây đau lại là những lỗi kéo luôn process đi. Riêng trường hợp timeout còn tệ hơn lỗi thường: bạn không phân biệt được "broker chưa hề nhận" với "broker nhận rồi nhưng mất ack". Retry ở tình huống đầu thì đúng; retry ở tình huống sau là bạn vừa publish hai lần.
Đảo lại, publish trước rồi mới commit, chỉ đổi chiều thiệt hại — giờ bạn có thể phát sự kiện cho một đơn hàng đã bị rollback, khó phát hiện hơn và khó gỡ hơn nhiều.
Không có thứ tự nào giữa hai lần ghi độc lập khiến chúng trở thành nguyên tử. Vấn đề nằm ở hệ thống thứ hai, không phải ở thứ tự.
Outbox
Vậy thì đừng ghi vào hệ thống thứ hai nữa.
Sự kiện trở thành một dòng trong chính database đó, ghi trong cùng transaction với thay đổi nghiệp vụ:
BEGIN;
INSERT INTO orders (id, customer_id, total) VALUES (…);
INSERT INTO outbox (aggregate_id, type, payload)
VALUES (…, 'order.placed', '{"orderId": …}');
COMMIT;
Hoặc cả hai dòng cùng tồn tại, hoặc không dòng nào cả. Ý tưởng chỉ có vậy, và đó cũng là phần duy nhất bắt buộc phải chính xác tuyệt đối.
Một relay riêng đọc các dòng chưa publish rồi đẩy sang Kafka. Nếu nó chết giữa chừng, các dòng vẫn còn nguyên và nó chạy tiếp từ chỗ dang dở. Nếu nó publish xong rồi chết trước khi đánh dấu đã gửi, khi khởi động lại nó sẽ publish dòng đó lần nữa.
Câu cuối không phải lỗi. Đó chính là cam kết: giao ít nhất một lần (at-least-once).
Publish mà không nói dối về thứ tự
Hai chi tiết quyết định pattern này có trụ được dưới tải hay không.
Chia partition theo aggregate, đừng chia vòng tròn. Mọi sự kiện của một đơn hàng phải vào cùng một partition Kafka, lấy order id làm key. Kafka chỉ bảo đảm thứ tự trong một partition, và không bảo đảm gì giữa các partition. Lấy aggregate làm key thì order.placed không bao giờ đến sau order.cancelled của cùng đơn. Lấy thứ khác làm key thì sớm muộn cũng xảy ra.
Giành dòng trước khi gửi. Sẽ có nhiều instance relay chạy song song. SELECT … FOR UPDATE SKIP LOCKED cho mỗi instance lấy một lô riêng mà không chặn nhau:
SELECT id, type, payload
FROM outbox
WHERE published_at IS NULL
ORDER BY id
FOR UPDATE SKIP LOCKED
LIMIT 100;
SKIP LOCKED chính là thứ làm relay scale ngang được. Thiếu nó, instance thứ hai ngồi chờ instance thứ nhất, và bạn có một cách rất tốn kém để chạy đúng một worker.
Nửa còn lại ở phía consumer: inbox
At-least-once nghĩa là trùng lặp là chuyện bình thường, nên consumer phải nhìn thấy cùng một sự kiện hai lần mà không làm việc hai lần.
Mỗi consumer giữ một bảng inbox với ràng buộc unique trên event id. Lệnh insert nằm cùng transaction với tác động nghiệp vụ:
await tx.inbox.create({ data: { eventId } }) // lần thứ hai sẽ ném lỗi
await tx.loyalty.increment({ customerId, points })
Lần giao thứ hai vi phạm unique, transaction rollback, và không có gì bị cộng đôi. Giao ít nhất một lần, nhưng tác động đúng một lần.
Đây là chỗ các nhóm hay cắt bớt, và cũng là nửa quyết định cam kết kia có thật hay không. Outbox mà không có inbox chỉ đẩy bài toán trùng lặp xuống service kế tiếp.
Cái giá phải trả
Nó không miễn phí, và giả vờ ngược lại là cách người ta lãnh đủ về sau:
- Độ trễ. Polling thêm độ trễ. Chu kỳ một giây thì người dùng không cảm nhận được và database chịu được thoải mái; xuống dưới một giây thì không đáng. Change data capture bỏ hẳn polling, đổi lại phải vận hành Debezium.
- Một bảng phình ra. Outbox cần dọn. Các dòng đã publish quá vài ngày thì xoá theo lịch — giữ đủ để debug, không giữ đủ để thành gánh nặng.
- Thứ tự chỉ đúng trong phạm vi aggregate. Thứ gì cần thứ tự toàn cục thì cần thiết kế khác. Thực tế thì không có thứ nào cần.
- Payload là một hợp đồng. Dòng ghi tuần trước sẽ được đọc bởi code deploy hôm nay, nên payload phải có version và consumer phải bỏ qua field lạ.
Nếu làm lại, tôi sẽ khác
Tôi sẽ viết inbox trước. Chúng tôi dựng outbox, thấy bản ghi trùng chạy về, rồi mới thêm khử trùng lặp trong lúc gấp deadline. Hai nửa này là một pattern, và nửa sau mới là thứ làm nửa trước trở thành sự thật.
Tôi cũng sẽ cưỡng lại ý định nhét cả snapshot aggregate vào payload. Nghe rất hấp dẫn — consumer khỏi phải gọi lại — nhưng mỗi field khi đó thành một hợp đồng không sửa được. Chỉ gửi định danh và đúng những dữ kiện đã thay đổi thì sống lâu hơn nhiều.
Mười ba service, một nguồn sự thật duy nhất, và không bao giờ ghi thẳng lên broker. Pattern chỉ có vậy.