Consumer Group: Điều gì xảy ra sau khi một Event được tạo ra?
Nguyễn Trung Dũng
Trong phần giới thiệu về EDA, chúng ta biết rằng mọi thứ bắt đầu từ Business Event, và các service sẽ không giao tiếp với nhau bằng cách gọi trực tiếp API. Vậy điều gì sẽ xảy ra sau khi một event OrderCreated được tạo ra? Làm thế nào để Inventory biết để giảm số lượng hàng còn lại trong kho, hay Analytics biết để cập nhật các metrics? Để tiếp nối, hôm nay chúng ta sẽ cùng tìm hiểu về Consumer Group, cơ chế giúp hiện thực hoá kiến trúc EDA, đồng thời mang lại khả năng mở rộng ấn tượng của Event Broker, với đại diện tiêu biểu là Kafka.
Nhắc lại kiến trúc truyền thống, nếu Order service muốn thông báo cho Inventory hay Analytics, nó thường phải gọi API của từng service, điều này gây ra một số vấn đề:
Thứ nhất là latency. Flow chính sẽ phải chờ các API trả về kết quả trước khi được coi là thực sự hoàn thành. Ngay cả khi các API calls được thực hiện song song, chỉ cần một service phản hồi chậm thì response time của toàn bộ request cũng bị kéo dài.
Thứ hai là maintainability. Mỗi khi xuất hiện một service mới, Order lại phải biết thêm về endpoint, giao thức và cách giao tiếp với service đó. Theo thời gian, số lượng dependency ngày càng tăng, khiến hệ thống trở nên khó mở rộng và khó bảo trì.
Với EDA, mọi thứ diễn ra theo cách hoàn toàn khác.
Khi đơn hàng được tạo, Order service chỉ có nhiệm vụ duy nhất là tạo ra một event OrderCreated và đưa nó vào Event Broker. Từ thời điểm này, bất kì service nào quan tâm tới event đó chỉ cần consume từ Event Broker sau đó xử lý một cách độc lập. Inventory cập nhật hàng trong kho, Analytics ghi nhận số liệu... Mỗi service đảm nhận một trách nhiệm riêng, nhưng đều bắt đầu từ cùng một Business Event.
Một event có thể được nhiều service cùng xử lý. Event Broker có thể phân phối event đó cho Inventory / Notification / Analytics như thế nào? Ở Message Queue truyền thống, mỗi khi event được consume bởi một service, nó lập tức được đưa ra khỏi queue, dẫn đến việc một event chỉ được xử lý một lần bởi một consumer.
Là một Event Broker hiện đại, Kafka tiếp cận theo hướng khác. Thay vì xoá event sau khi được consume, Kafka lưu trữ event trong Topic (có thể hình dung đơn giản là một dạng queue) trong một khoảng thời gian nhất định (mặc định 7d).
Trong khi Producer vẫn liên tục gửi event vào topic nếu có, Consumer sẽ tự theo dõi vị trí mà mình đã xử lý thông qua chỉ số gọi là "Offset". Sau khi xử lý xong một event, Consumer sẽ cập nhật offset lên Broker để đánh dấu tiến trình. Nhờ cơ chế này, Consumer có thể được stop, update, restart, migrate... mà không sợ mất event, thậm chí đọc lại các event cũ bằng cách đẩy offset về phía đầu của queue.
Đến đây, chúng ta có thể thấy Producer và Consumer hoàn toàn độc lập. Tuy nhiên cần chú ý rằng tạo ra một event thường nhanh, nhưng để xử lý xong một event thì mất nhiều thời gian hơn rất nhiều.
Khi có quá nhiều event đến cùng một lúc (những đợt sale lớn), lượng event chưa được xử lý ngày càng tăng. Nếu tình trạng này kéo dài, chúng ta có thể thấy vấn đề về delay: đã xong order, thanh toán xong, nhưng 30' sau mới có confirmation email. Về mặt logic, đơn hàng vẫn được xử lý dù bị chậm, nhưng về mặt sản phẩm, trải nghiệm rõ ràng bị ảnh hưởng.
Chúng ta đi đến bài toán tiếp theo:
Khi có quá nhiều event, Chuyện gì xảy ra và solution là gì cho việc tốc độ xử lý của Consumer không theo kịp tốc độ tạo ra event của Producer?
Câu trả lời là cơ chế xử lý song song, bằng cách dùng nhiều Consumer hơn.
Nhưng nếu chỉ đơn giản tăng thêm các Consumer mà không có cơ chế phối hợp, những Consumer sẽ đọc lại toàn bộ event trong topic, vì chúng không biết các event này đã được xử lý chưa.
Kafka giải quyết bài toán này bằng cơ chế Consumer Group.

Consumer Group là tập hợp các Consumer cùng chia sẻ công việc xử lý một Topic. Thay vì để một Consumer xử lý toàn bộ event, các Consumer trong cùng một Group sẽ cùng nhau chia sẻ workload, đảm bảo mỗi event chỉ được xử lý bởi một Consumer.
Khả năng xử lý song song chính là giá trị quan trọng nhất của Consumer Group.
Về mặt kỹ thuật, Kafka thực hiện cơ chế này bằng cách chia một Topic thành nhiều Partition, sau đó phân chia các Partition này cho các Consumer trong Group. Với mỗi Partition, một Consumer Group chỉ có một Offset tại một thời điểm, và mỗi Partition cũng chỉ được assign cho một Consumer trong Group. Nhờ đó, các Consumer có thể cùng xử lý nhiều event song song mà không bị trùng.
Cơ chế này cho thấy số lượng Partition là một yếu tố quan trọng quyết định khả năng xử lý song song của một Consumer Group. Nếu muốn scale Consumer Group, chúng ta không thể chỉ tăng số lượng Consumer mà còn phải có đủ Partition để các Consumer có thể thực sự chia sẻ workload.

Hãy cùng quay lại với câu chuyện về Sale với rất nhiều event được Order service đưa vào Event Broker. Từ đây, các service khác nhau có thể sử dụng event đó cho các mục đích khác nhau:
Inventory sử dụng để cập nhật tồn kho
Analytics dùng để ghi nhận số liệu phân tích
Notification dùng để gửi email xác nhận
Mỗi workload sẽ sử dụng một Consumer Group riêng, nhờ vậy các event có thể được xử lý hoàn toàn độc lập, mỗi workload có thể scale theo nhu cầu riêng. Ví dụ, khi lượng đơn hàng tăng mạnh, Notification service sẽ tăng số lượng consumer để theo kịp lượng lớn event mà không ảnh hưởng gì tới Inventory hay Analytics.
Điều quan trọng tiếp theo là Order service hoàn toàn không cần biết có những service nào đang xử lý event. Một Consumer Group cho một service mới có thể được thêm vào mà không cần thay đổi bất cứ component đang chạy nào.
Consumer Group là một cơ chế quan trọng giúp Event-Driven Architecture có thể mở rộng linh hoạt mà không làm ảnh hưởng tới các hệ thống khác, nhờ đặc tính loose coupling giữa các service.
Nhưng khi một Event đã trở thành trung tâm giao tiếp giữa các service, một câu hỏi quan trọng khác xuất hiện: điều gì xảy ra nếu Event bị mất, được xử lý nhiều lần, hoặc Consumer gặp sự cố giữa chừng? Trong hệ thống phân tán, việc gửi được một Event mới chỉ là bước đầu. Đảm bảo Event được xử lý đúng mới thực sự là bài toán khó.
Đó cũng sẽ là nội dung chúng ta tìm hiểu trong bài viết tiếp theo.