# Gắn AI vào hàng đợi sự kiện sẵn có của khách mà không phải viết lại hệ thống

> Khách đã có hàng đợi chạy ổn định. Việc của FDE là gắn một consumer AI vào đó mà không làm thanh toán chạy hai lần, không để lỗi dồn lại âm thầm và không tạo ra cơn bão sự kiện.

Bản gốc: https://fdetimes.net/vi/bach-khoa/event-driven-event-sourcing-cqrs/

Tuần đầu ở chỗ khách, bạn được giao một yêu cầu nghe khá đơn giản: dùng LLM tự động gắn nhãn cho mỗi ticket hỗ trợ ngay khi ticket được tạo. Nhóm kỹ thuật của khách mở sơ đồ hệ thống. Mỗi khi người dùng gửi ticket, giao diện đẩy một message vào hàng đợi, rồi vài dịch vụ phía sau lấy message đó ra để xử lý.

Lúc này nhiều kỹ sư sẽ bắt đầu vẽ lại cả kiến trúc: thêm event store, tách read model, áp CQRS cho bài bản. Đó thường là bước đi sai. Khách không thuê bạn để làm lại hệ thống. Họ cần một consumer mới chạy trên luồng sự kiện họ đã có, chạy an toàn và không làm hỏng thứ đang hoạt động.

Để làm tốt việc này, bạn cần phân biệt rõ ba khái niệm hay bị dùng lẫn với nhau, sau đó viết một consumer tránh được bốn cái bẫy quen thuộc.

## Ba khái niệm, ba mức cam kết khác nhau

Event-driven là mức đơn giản nhất. Theo tài liệu kiến trúc Azure của Microsoft, một background job có thể được khởi động khi giao diện hoặc một job khác đặt message vào hàng đợi. Message đó chứa dữ liệu về một hành động vừa xảy ra, ví dụ người dùng vừa đặt hàng. Hệ thống ticket ở trên thuộc đúng kiểu này.

Event sourcing đòi hỏi nhiều hơn hẳn. Martin Fowler định nghĩa nó là cách ghi lại mọi thay đổi của trạng thái ứng dụng dưới dạng một chuỗi sự kiện. Microsoft mô tả cụ thể hơn: toàn bộ chuỗi hành động được lưu trong một kho chỉ cho phép ghi thêm (append-only), và kho này là nguồn dữ liệu gốc thay cho bảng trạng thái hiện tại.

CQRS tách thao tác ghi (command) và thao tác đọc (query) thành hai mô hình dữ liệu riêng để tối ưu độc lập. Khi hai bên dùng hai kho khác nhau, cách phổ biến là để mô hình ghi phát sự kiện mỗi khi cập nhật database, và mô hình đọc dựa vào đó để làm mới dữ liệu.

Còn một nhầm lẫn nữa cần gỡ. Microsoft lưu ý rằng không nên coi event store và message broker là một. Kafka hay một hàng đợi tương tự chỉ là lớp phân phối, giúp đưa sự kiện tới các projection và consumer bên ngoài. Nếu khách bảo "chúng tôi có Kafka rồi", điều đó chưa có nghĩa là họ đang dùng event sourcing.

## Vì sao không nên đề xuất tái kiến trúc

Chính Microsoft cảnh báo rằng event sourcing tốn kém và hiếm khi cần đến. Với phần lớn hệ thống và phần lớn các thành phần trong một hệ thống, cách quản lý dữ liệu truyền thống là đủ.

Fowler cũng nói tương tự về CQRS: đây là một bước nhảy lớn về tư duy cho tất cả mọi người liên quan, nên chỉ làm khi lợi ích thực sự xứng đáng.

Khi bàn về event sourcing, Microsoft nêu một ưu điểm: code phát sự kiện tách rời khỏi các hệ thống đăng ký nhận sự kiện. Lý lẽ tương tự vẫn áp dụng được cho một hàng đợi thông thường, dù không có event sourcing: giao diện chỉ đặt message vào hàng đợi, không cần biết phía sau có bao nhiêu job đang lấy message ra.

Vì vậy bạn có thể thêm một consumer AI mà không phải sửa một dòng code nào ở ứng dụng đang phát sự kiện. Với FDE, điều này rất có lợi: phạm vi nhỏ, ít đụng chạm nội bộ và có thể rollback chỉ bằng cách tắt consumer.

**Điểm mấu chốt:** Gắn AI vào luồng sự kiện khách đang có, đừng đòi khách xây lại luồng đó.

## Ví dụ: consumer gắn nhãn ticket

Giả sử mỗi message có dạng `TicketCreated` với `event_id`, `ticket_id`, `correlation_id` và nội dung ticket. Consumer dưới đây là phiên bản bạn có thể triển khai thật, không chỉ để demo:

```python
def handle(msg):
try:
event = parse_ticket_created(msg.body)
except SchemaError:
msg.dead_letter(reason="schema")   # không retry message hỏng
return

if processed_events.exists(event.event_id):
msg.complete()                     # message trùng: bỏ qua
return

if labels.exists_for_correlation(event.ticket_id, event.correlation_id):
msg.complete()                     # vòng lặp: ticket đã có nhãn trong cùng nghiệp vụ
return

label = classify_with_llm(event.text)  # có thể chạy lại, không gây hại

with db.transaction():
labels.upsert(event.ticket_id, label,
source_event=event.event_id,
correlation_id=event.correlation_id)
processed_events.insert(event.event_id)
outbox.add("TicketLabeled",
ticket_id=event.ticket_id,
correlation_id=event.correlation_id,
produced_by="ai-labeler")
msg.complete()
```

Đọc từng khối từ trên xuống, bạn sẽ thấy mỗi khối xử lý một bẫy riêng. Phần quan trọng nhất là ba thao tác trong `db.transaction()`: lưu nhãn, đánh dấu sự kiện đã xử lý và ghi vào bảng outbox.

Ba việc này cùng thành công hoặc cùng thất bại. Microsoft khuyến nghị kết hợp Transactional Outbox với consumer idempotent để giữ mô hình ghi và mô hình đọc nhất quán, và nguyên tắc đó áp dụng y nguyên cho consumer AI ở đây.

Hãy để ý rằng lời gọi LLM nằm ngoài transaction. Nếu tiến trình chết sau khi gọi model nhưng trước khi commit, message sẽ được giao lại và model bị gọi thêm một lần. Bạn mất thêm một ít chi phí token, nhưng không có tác dụng phụ nào bị lặp lại. Đây là một đánh đổi chấp nhận được.

## Bốn cái bẫy và cách đoạn code tránh chúng

Bẫy đầu tiên là message bị giao trùng. Microsoft nói rõ hàng đợi chỉ đảm bảo giao ít nhất một lần, tức là cùng một message có thể đến nhiều lần.

Nếu handler không idempotent, projection sẽ dần lệch khỏi luồng sự kiện, còn các tác dụng phụ như thanh toán hay gửi thông báo có thể chạy hơn một lần. Với một consumer AI có quyền gửi email cho khách hàng, điều này có thể thành sự cố ngay trong tuần đầu.

Bẫy thứ hai là message hỏng (poison message). Có thể một ticket có nội dung rỗng, hoặc dùng schema cũ khiến parser báo lỗi. Microsoft nhắc rằng mỗi message nằm trong dead-letter queue là một phần việc chưa hoàn thành. Nếu không ai theo dõi hàng đợi này, lỗi sẽ dồn lại mà không ai biết, và khách chỉ phát hiện khi thấy nhiều ticket không có nhãn.

Bẫy thứ ba là event storm. Mô hình choreography cho phép mỗi dịch vụ tự quyết định khi nào và xử lý thế nào, thường thông qua broker theo kiểu publish-subscribe. Microsoft cảnh báo kiểu thiết kế này có thể vô tình tạo ra vòng lặp phản hồi hoặc event storm.

Thử hình dung một dịch vụ khác nghe `TicketLabeled` và cập nhật ticket, việc cập nhật đó phát lại `TicketCreated`, và consumer AI lại chạy thêm một lần nữa. Sự kiện mới này do dịch vụ kia phát ra, nên tự nó sẽ không mang `produced_by="ai-labeler"`; kiểm tra `event_id` cũng không chặn được vì đây là một sự kiện mới.

Cách chặn là thỏa thuận với nhóm của khách để mọi dịch vụ phía sau chuyển tiếp `correlation_id` sang sự kiện chúng phát ra.

Khi đó bước `labels.exists_for_correlation` trong đoạn code mới phát huy tác dụng: ticket đã có nhãn trong cùng `correlation_id` thì bị bỏ qua, và vòng lặp dừng ngay ở lượt thứ hai.

Thiếu phần chuyển tiếp từ phía khách, bước kiểm tra này không có gì để so.

Bẫy thứ tư là dữ liệu cũ. Khi kho đọc và kho ghi tách riêng, hệ thống chỉ nhất quán sau một độ trễ (eventual consistency), nên dữ liệu đọc có thể chưa phản ánh thay đổi mới nhất.

Nếu agent của bạn đọc lịch sử khách hàng từ read model để quyết định hoàn tiền, nó có thể quyết định dựa trên dữ liệu cũ vài giây. Nếu quyết định cần dữ liệu mới nhất, hãy lấy thông tin trực tiếp từ payload của sự kiện, hoặc hỏi nhóm của khách xem read model thường trễ bao lâu.

## Correlation ID: để còn trả lời được "chuyện gì đã xảy ra?"

Khi không có bộ điều phối trung tâm, Microsoft lưu ý rằng không thành phần nào thấy được toàn bộ một nghiệp vụ đang diễn ra. Vì thế cần có correlation ID và distributed tracing. Khi một khách hàng phàn nàn rằng ticket của họ bị gắn sai nhãn, bạn cần lần ngược từ nhãn đó về sự kiện gốc, rồi tới lời gọi model và prompt đã dùng.

Cách làm rất đơn giản. Mỗi sự kiện bạn phát ra phải mang theo `correlation_id` của sự kiện đã kích hoạt nó. Mỗi dòng log của consumer cũng ghi kèm ID này. Nếu khách chưa có ID như vậy, hãy đề xuất thêm nó trước khi đưa AI lên production. Đây là thay đổi nhỏ, nhưng rất có ích cho buổi điều tra sự cố đầu tiên.

## Trình tự làm tại chỗ khách

Trước hết, xác định khách đang ở mức nào trong ba mức trên. Microsoft cho rằng event sourcing phù hợp nhất khi ứng dụng vốn đã dùng sự kiện như một phần tự nhiên trong cách vận hành.

Nhiều khách sẽ chỉ có hàng đợi, và như vậy là đủ. Tiếp theo, xin một mẫu message thật, đọc schema và hỏi họ đã từng đổi schema bao giờ chưa.

Sau đó, chạy consumer ở chế độ shadow: vẫn gắn nhãn nhưng chỉ ghi vào một bảng riêng, chưa phát sự kiện nào ra ngoài. Khi số message trong dead-letter queue ổn định và nhãn đạt yêu cầu, bạn mới bật outbox. Cách làm này giúp bạn tách hai nỗi lo: chất lượng của model và độ an toàn của hệ thống.

Nếu đang chuẩn bị chuyển sang vai trò FDE, bạn nên đưa đúng những chi tiết này vào CV. Câu "Tích hợp LLM vào luồng Kafka" không nói lên nhiều điều.

Câu "Viết consumer idempotent bằng outbox, có cảnh báo dead-letter và correlation ID, chạy shadow hai tuần trước khi phát sự kiện" cho nhà tuyển dụng thấy bạn hiểu rủi ro khi làm việc trên hệ thống của người khác.

Khi một JD yêu cầu tích hợp với hệ thống sẵn có hoặc làm việc với message queue, hãy chuẩn bị kể lại đúng một ví dụ như trên trong buổi phỏng vấn, kèm cách bạn xử lý message bị giao trùng.

Phần dễ gây ấn tượng trong một buổi demo là model gắn nhãn chính xác. Còn thứ quyết định khách có tiếp tục tin bạn hay không là việc consumer vẫn an toàn khi cùng một message được giao tới lần thứ hai.

**Thử ngay tuần này:**

- Chọn một hàng đợi bạn đang có quyền truy cập, gửi cùng một message hai lần và kiểm tra xem consumer có tạo ra hai kết quả hay không.
- Viết lại consumer theo mẫu trong bài: kiểm tra event_id đã xử lý, bỏ qua ticket đã có nhãn trong cùng correlation_id, ghi kết quả và outbox trong cùng một transaction.
- Mở dead-letter queue của một hệ thống bạn quản lý và đếm số message đang nằm trong đó. Nếu chưa có cảnh báo cho con số này, hãy thêm một cảnh báo.

## Nguồn

- [Best Practices for Background Jobs - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/best-practices/background-jobs)

- [Event Sourcing Pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/event-sourcing)

- [CQRS Pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/cqrs)

- [Choreography Pattern - Azure Architecture Center | Microsoft Learn](https://learn.microsoft.com/en-us/azure/architecture/patterns/choreography)

- [Event Sourcing (Martin Fowler)](https://martinfowler.com/eaaDev/EventSourcing.html)

- [CQRS (Martin Fowler)](https://martinfowler.com/bliki/CQRS.html)
