Tách làm hai tầng.
Tầng thuần logic (phần lớn test): hàm xử lý nhận vào một message đã parse và trả về kết quả, không tự kết nối broker. Test ở đây nhanh và bao được các case quan trọng:
- Message hợp lệ → tạo đúng bản ghi, phát đúng sự kiện tiếp theo.
- Idempotency: cùng một message xử lý hai lần chỉ có tác dụng một lần. Đây là case bắt buộc vì broker đảm bảo at-least-once, message lặp là chuyện bình thường chứ không phải bất thường.
- Message hỏng: thiếu field, sai kiểu, phiên bản schema cũ → không được ném lỗi vô hạn khiến consumer kẹt; phải đẩy sang dead-letter queue.
- Lỗi tạm thời (DB mất kết nối) → không ack, để message quay lại và retry.
await handleOrderPaid(msg)
await handleOrderPaid(msg) // duplicate delivery
expect(await countPayments(msg.orderId)).toBe(1)Tầng tích hợp (ít test, chạy chậm hơn): dựng broker thật bằng Testcontainers, publish một message, chờ tới khi hệ quả xuất hiện. Ở đây kiểm những thứ chỉ broker thật mới lộ: serialize/deserialize, cấu hình consumer group, ack/nack, định tuyến sang DLQ.
Lưu ý về flaky: xử lý là bất đồng bộ, nên đừng dùng sleep cố định. Dùng vòng lặp chờ có timeout ("poll cho tới khi bản ghi xuất hiện, tối đa 5 giây") — vừa nhanh hơn khi máy khoẻ, vừa không đỏ ngẫu nhiên khi CI chậm.