Ba mươi tư bài vừa rồi nói về hàng đợi: thông điệp vào, consumer lấy ra, thông điệp biến mất. Stream đảo ngược điều cuối cùng — và vì thế nó là một mô hình khác, không phải một tối ưu hoá khác.
Ba kiểu, ba mô hình
| Classic | Quorum | Stream | |
|---|---|---|---|
| Thông điệp sau khi ack | biến mất | biến mất | còn nguyên |
| Hai consumer trên cùng một cái | chia việc | chia việc | mỗi bên đọc đủ |
| Đọc lại từ đầu | không | không | được |
| Nhân bản qua cụm | không | có | có |
Thông lượng trên cùng cụm ba node
Cùng phép đo với phần trước — 50 000 thông điệp 200 byte, persistent, có confirms:
| Kiểu | Thông lượng |
|---|---|
| classic | 130 330 msg/s |
| stream | 79 025 msg/s |
| quorum | 44 043 msg/s |
Stream nhanh gấp 1,8 lần quorum trong khi vẫn có bản sao trên cả ba node. Lý do nằm ở cấu trúc: nó chỉ ghi thêm vào cuối một tệp log, không phải quản lý trạng thái từng thông điệp như hàng đợi.
Hai consumer cùng đọc đủ
1 000 thông điệp, 2 consumer độc lập -> A=1000, B=1000, TỔNG=2000
So với phần 9, nơi hai consumer trên một hàng đợi chia nhau 50/50 và phải dựng fanout với hai hàng đợi riêng mới có được 100/100. Stream cho kết quả đó mà không cần exchange, không cần binding, không cần nhân bản dữ liệu — một bản trên đĩa, nhiều người đọc.
Đây là khác biệt kiến trúc thật sự: với hàng đợi, số bên nghe quyết định số bản ghi (phần 9 đo được fanout 100 hàng đợi làm thông lượng tụt 68 lần). Với stream, thêm bên nghe gần như miễn phí về phía ghi.
Đọc lại từ đầu
consumer thứ ba đọc lại từ đầu -> 1 000 thông điệp
x-stream-offset: first cho một consumer mới đọc lại toàn bộ lịch sử còn giữ. Với hàng đợi thì điều này là không thể — thông điệp đã ack là đã đi.
Đây là thứ mở ra những việc hàng đợi không làm được: dựng lại trạng thái của một dịch vụ mới từ lịch sử sự kiện, chạy lại một ngày dữ liệu sau khi sửa lỗi, cho một đội khác đọc cùng luồng mà không ảnh hưởng ai.
Cái bẫy: số thông điệp luôn báo 0
Ghi 1 000 thông điệp vào một stream rồi hỏi cả ba nguồn:
queueDeclarePassive().getMessageCount() = 0
rabbitmqctl list_queues : sx stream 0
HTTP API : messages = 0, message_bytes = null
Ba nguồn, cùng một số 0 — trong khi consumer thứ ba vẫn đọc ra đủ 1 000 thông điệp.
Theo dõi stream phải theo độ trễ của từng consumer (offset nó đang đọc so với offset cuối), không phải theo độ sâu. Và giới hạn dung lượng phải đặt bằng x-max-length-bytes hoặc thời gian giữ, vì không có cơ chế "consumer rút hết thì hàng đợi rỗng" nào cả.
Vài luật cứng
basicQos là bắt buộc:
consumer stream không đặt prefetch
-> 406 PRECONDITION_FAILED - consumer prefetch count is not set for stream queue 'sq'
Broker từ chối thẳng. Hợp lý: không có trần thì consumer sẽ bị dội cả lịch sử vào mặt.
Không có auto-delete, không có x-max-priority — cả hai đều 406. Stream là hạ tầng dữ liệu, không phải hàng đợi tạm.
Khi nào stream, khi nào là dấu hiệu cần Kafka
Dùng stream khi nhiều bên cần đọc cùng một luồng sự kiện một cách độc lập, hoặc khi bạn cần đọc lại lịch sử. Nếu bạn đang định dựng một fanout với mười hàng đợi để mười dịch vụ cùng nghe, stream làm việc đó rẻ hơn nhiều.
Đừng dùng stream khi công việc là hàng đợi việc — mỗi thông điệp một người làm, làm xong thì biến mất. Stream không có basicReject, không có DLX, không có ưu tiên; bạn sẽ phải tự viết lại tất cả những thứ mà 34 bài trước đo được là RabbitMQ đã làm sẵn.
Đó là dấu hiệu cần Kafka khi bạn cần giữ lịch sử hàng tháng với hàng chục terabyte, phân vùng theo khoá để xử lý song song có thứ tự, hoặc một hệ sinh thái xử lý luồng đầy đủ. Stream của RabbitMQ giải quyết tốt trường hợp "tôi cần đọc lại và nhiều bên cùng nghe" trong một hệ đã dùng RabbitMQ — nó không cố thay Kafka, và dùng nó như Kafka là chọn nhầm công cụ.
Bài sau: Shovel và Federation — nối hai broker với nhau.
Thử ba mươi giây
docker exec -u rabbitmq n1 rabbitmqctl -q list_queues name type members messages
Dòng nào type=stream mà messages=0 thì đừng vội mừng — hãy đi tìm offset của consumer thay vì tin vào cột cuối.