Nghĩ về một series phim trên dịch vụ xem trực tuyến. Xem xong một tập, tập đó không biến mất; mười người cùng xem, ai theo nhịp nấy; một người mới vào hôm nay bắt được từ tập 1. Và câu hỏi "còn mấy tập đang chờ bạn" là vô nghĩa — cái đáng hỏi là "bạn đang ở tập mấy so với tập mới nhất". Stream của RabbitMQ là đúng cái kho phim đó, khác hẳn hàng đợi kiểu băng chuyền nơi món đồ rơi khỏi cuối băng khi có người nhặt. Nó là một mô hình khác, không phải một tối ưu hoá khác — và bài này đo chỗ khác biệt đó.
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ụ.
Muốn kiểm nhanh, đừng tin cột messages:
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.
Mẫu số chung
Trước khi chọn công cụ nhắn tin nào, có một ngã ba mô hình quyết định gần hết mọi thứ về sau: đọc là tiêu huỷ, hay đọc là quan sát? Hàng đợi kiểu băng chuyền trao mỗi món cho đúng một người rồi món đó biến mất — hợp với việc cần làm. Log ghi-thêm giữ nguyên mọi thứ và mỗi người đọc mang một con trỏ riêng — hợp với sự kiện cần nhiều bên quan sát, và cần tua lại. Chọn nhánh nào là đã định luôn: fan-out đắt hay rẻ, có replay được không, giữ lịch sử ra sao. Cùng ngã ba ấy hiện ra ở SQS với Kafka, một message queue với một event log, một danh sách việc với một sổ kiểm toán, pop một phần tử với đọc theo chỉ số. Câu hỏi gọn: người đọc nên "lấy đi rồi thôi", hay "xem theo nhịp của mình mà không đụng người khác"?
Điều thứ hai, và nó bẫy ngay cả người có kinh nghiệm: một chỉ số chỉ có nghĩa trong đúng mô hình nó được định nghĩa. "Độ sâu hàng đợi" là thước đo sức khoẻ hoàn hảo cho hàng đợi, nhưng bê nguyên sang stream thì nó đọc ra 0 một cách vô nghĩa — không phải sai, mà là không áp dụng, và một cái đồng hồ đọc số vô nghĩa còn nguy hơn không có đồng hồ. Thước đo đúng cho log là độ trễ offset của từng consumer. Cùng cái lỗi bê-chỉ-số-qua-mô-hình ấy ở việc canh "queue depth" cho một hệ pub/sub, đo CPU% cho một việc nghẽn I/O, nhìn hit-rate của cache cho một kho ghi-thẳng. Khi đổi sang một mô hình mới, đừng mang theo bảng đồng hồ cũ — hỏi lại từ đầu: trong mô hình này, "khoẻ" nghĩa là con số nào?
Bài sau: Shovel và Federation — nối hai broker với nhau.