Hình dung một băng chuyền: hàng cứ trôi tới bàn bạn bất kể bạn có kịp nhặt hay không, và nếu bạn không đặt một cái chặn thì nó dồn thành đống đè lên bạn. Đó là mô hình đẩy (push). Ngược lại là quầy buffet: bạn tự lấy đĩa tiếp theo khi rảnh tay — mô hình kéo (pull), tự nó có nhịp. RabbitMQ đẩy bản tin xuống consumer, nên nếu bạn không tự đặt cái chặn, một hàng đợi đầy sẽ dội thẳng vào mặt ứng dụng. Bài này đo đúng cú dội đó — và nó nặng hơn nhiều người nghĩ.

Tôi đã có một sê-ri 40 bài về RabbitMQ đo gần như mọi thứ trong broker. Bài này hỏi một câu hẹp hơn: khi nối RabbitMQ từ Vert.x, những gì đã đo còn đúng không, và có gì mới hỏng?

Câu trả lời ngắn: các con số về broker vẫn đúng, còn cái mới hỏng thì hỏng rất nặng.

RabbitMQ 4.3.5 trong Docker, vertx-rabbitmq-client, bản tin 200 byte.

Gửi: bốn cách, chênh nhau 35 lần

Thông lượng Mỗi bản tin
Không confirm, bắn hết rồi chờ 197 126 bản tin/giây 5,1 µs
Không confirm, tuần tự 47 637 bản tin/giây 21,0 µs
Có confirm, tuần tự 5 649 bản tin/giây 177,0 µs
Có confirm, bắn hết rồi chờ 92 180 bản tin/giây 10,8 µs

Ba điều rút ra, và điều thứ nhất là cái bẫy lớn nhất:

"Không confirm" hoàn thành sau 21 µs vì nó không hỏi broker gì cả. Future mà basicPublish trả về hoàn thành khi khung tin đã được ghi ra socket — không phải khi broker nhận được, càng không phải khi bản tin đã nằm trên đĩa. Nếu bạn coi onSuccess là "đã gửi xong an toàn" thì bạn đang nhầm; đây đúng là chuyện publisher confirms trong sê-ri RabbitMQ đã bàn, chỉ khác là API Future làm nó trông như một lời hứa mạnh hơn thực tế.

Confirm tuần tự đắt 177 µs — đó mới là một vòng đi về thật tới broker và về.

Nhưng confirm không bắt buộc phải đắt. Vẫn dùng confirm mà không chờ từng cái, thông lượng lên 92 180 bản tin/giây — chỉ kém bản không confirm hai lần, và gấp 16 lần bản confirm tuần tự. Bạn giữ được bảo đảm mà không trả 177 µs mỗi bản tin, miễn là đừng biến mỗi lần gửi thành một điểm đồng bộ.

Chiều nhận, với autoAck, cho 133 230 bản tin/giây.

Chỗ Vert.x làm mọi thứ khác đi

Đây là phần không có trong sê-ri RabbitMQ, vì nó không phải chuyện của broker.

Tôi dựng một Verticle có cả HTTP server và consumer RabbitMQ — chuyện bình thường. Hàng đợi có sẵn 100 000 bản tin, handler làm một việc tốn CPU, autoAck bật. Rồi tôi curl đường HTTP:

truoc: 0,018 s  0,00073 s  0,00076 s  0,00069 s  0,00066 s
sau  : QUA-HAN(20 s)  QUA-HAN(20 s)  QUA-HAN(20 s)  10,19 s  0,00064 s

HTTP không chậm đi — nó chết hẳn hơn 50 giây, rồi sống lại đúng lúc bản tin cuối cùng được xử lý xong. Ba lần curl liên tiếp hết hạn 20 giây mà không nhận được byte nào.

Nguyên nhân là sự kết hợp của hai thứ, và cả hai đều là mặc định:

  • Handler của consumer chạy trên event loop, giống hệt chuyện với Kafka ở phần 36.
  • RabbitMQ đẩy bản tin xuống consumer, và với autoAck cùng hàng đợi nội bộ lớn thì không có gì kìm nó lại. Broker bơm cả 100 000 bản tin xuống nhanh hết mức, mỗi bản tin lại chiếm event loop thêm một nhịp.

Khác biệt với Kafka đáng chú ý: Kafka kéo theo lô nên max.poll.records tự tạo ra một nhịp nghỉ, và HTTP tuy tụt xuống 281 req/s nhưng vẫn trả lời. RabbitMQ đẩy, nên không có nhịp nghỉ nào và HTTP im lặng hoàn toàn. Cùng một lỗi lập trình, hậu quả nặng hơn hẳn.

Bản vá gồm hai nửa, thiếu nửa nào cũng không đủ:

rmq.start()
   .compose(x -> rmq.basicQos(10))                         // 1. gioi han so ban tin dang bay
   .compose(x -> rmq.basicConsumer("q-do", new QueueOptions()
        .setAutoAck(false)                                 //    autoAck=true lam basicQos vo nghia
        .setMaxInternalQueueSize(10)))
   .onSuccess(tt -> tt.handler(msg ->
        tho.executeBlocking(() -> xuLyNang(msg), false)    // 2. day viec ra khoi event loop
           .onComplete(ar -> rmq.basicAck(msg.envelope().getDeliveryTag(), false))));

Đo lại, vẫn hàng đợi đầy, vẫn cùng phép xử lý nặng:

sau  : 0,00099 s  0,00071 s  0,00065 s  0,00068 s  0,00065 s  0,00060 s

HTTP không hề hấn gì trong khi consumer vẫn chạy (đã xử lý 71 963 bản tin, còn 26 213 trong hàng đợi). basicQos(10) chính là prefetch mà sê-ri RabbitMQ đã đo — ở đây nó không chỉ ảnh hưởng thông lượng mà quyết định ứng dụng còn sống hay không.

Điểm dễ sai nhất: autoAck = true làm basicQos mất tác dụng. Prefetch đếm số bản tin chưa được xác nhận; tự động xác nhận ngay khi giao thì con số đó luôn bằng không và broker cứ đẩy tiếp. Hai tuỳ chọn này phải đi cùng nhau.

Hai bất ngờ từ RabbitMQ 4.x

Hàng đợi transient không độc quyền đã bị bỏ. Dòng khai báo quen thuộc queueDeclare(ten, false, false, true) — không bền, không độc quyền, tự xoá — làm broker đóng kết nối:

INTERNAL_ERROR - Feature `transient_nonexcl_queues` is deprecated

Lỗi trả về ở mức kết nối chứ không phải mức kênh, nên nó biểu hiện thành IOException kèm ShutdownSignalException và client tự nối lại rồi hỏng lại. Nếu bạn đang nâng cấp lên RabbitMQ 4.x thì đây là thứ nên tìm trước trong mã: đổi durable thành true.

AMQP 1.0 chạy sẵn, không cần plugin. rabbitmq-plugins list cho thấy rabbitmq_amqp1_0 vẫn ở trạng thái tắt:

[  ] rabbitmq_amqp1_0    4.3.5

Vậy mà vertx-amqp-client nối thẳng vào cổng 5672 và bắt tay thành công — RabbitMQ 4.x hỗ trợ AMQP 1.0 ngay trong lõi, plugin cũ chỉ còn là di tích. Nghĩa là bạn có hai client Vert.x để chọn:

vertx-rabbitmq-client vertx-amqp-client
Giao thức AMQP 0-9-1 AMQP 1.0
Khái niệm exchange, binding, queue — đúng mô hình RabbitMQ address, link — mô hình chung
Dùng khi chỉ nói chuyện với RabbitMQ cần đổi broker, hoặc nói với nhiều loại broker

Chọn cái đầu nếu bạn muốn dùng những thứ đặc trưng của RabbitMQ — exchange kiểu topic, DLX, quorum queue. Chọn cái sau nếu tính di động quan trọng hơn.

Muốn biết consumer của mình có phanh chưa, hỏi thẳng broker — không cần sửa mã:

docker exec <broker> rabbitmqctl list_consumers queue_name prefetch_count ack_required

prefetch_count bằng 0 nghĩa là không giới hạn — broker sẽ đẩy nhanh hết mức nó có thể, và nếu handler của bạn chạy trên event loop thì bạn đang cách một đợt dồn hàng đúng một bước. Cột ack_required bằng false thì prefetch có đặt cũng vô nghĩa.

Mẫu số chung

Khác biệt cốt lõi giữa hai cú "làm ngập consumer" — RabbitMQ chết hẳn còn Kafka chỉ chậm lại — nằm ở đẩy so với kéo. Một nguồn đẩy không có phanh sẵn: nó bơm nhanh hết mức nó có thể, và nếu bạn không tự dựng một cái chặn (prefetch, credit, cửa sổ) thì băng chuyền dồn hàng đè chết bạn. Một nguồn kéo thì cái phanh nằm sẵn trong thiết kế — bạn chỉ lấy khi gọi, nên max.poll.records tự giới hạn mỗi nhịp. Nguyên tắc: với mọi nguồn đẩy, việc kiểm soát lưu lượng là trách nhiệm của bạn, không phải của nguồn — một WebSocket server dội event, một callback đăng ký nhận thông báo, một stream phản ứng không có request(n), tất cả đều cần bạn tự đặt trần. Còn kéo thì an toàn hơn theo mặc định, đổi lại bạn phải chủ động đi lấy. Trước khi cắm một nguồn dữ liệu, hỏi: nó đẩy hay tôi kéo — và nếu nó đẩy, cái phanh của tôi đâu?

Điều thứ hai, một cái bẫy cấu hình: một tuỳ chọn an toàn có thể bị một tuỳ chọn hàng xóm âm thầm vô hiệu hoá. basicQos(10) là cái phanh — nhưng autoAck=true làm nó thành số 0, vì prefetch đếm bản tin chưa ack, mà tự-ack thì chẳng có cái nào chưa ack cả. Đặt đúng một nửa còn nguy hơn không đặt gì, vì bạn tưởng đã có phanh. Cùng cái bẫy "hai núm phải đi cùng nhau" ở khắp nơi: fsync vô nghĩa nếu ổ đĩa vẫn bật write-cache, một mức cô lập giao dịch không có hiệu lực nếu thiếu đúng kiểu khoá, một timeout kết nối không cứu được gì nếu thiếu timeout đọc. Bài học: một bảo đảm hiếm khi đến từ một cờ đơn lẻ — nó đến từ một bộ cờ nhất quán, nên khi bật một tính năng an toàn, hãy kiểm luôn cái hàng xóm có đang lặng lẽ tắt nó đi không.

Phần sau nối Redis và đo chênh lệch giữa gửi từng lệnh với gửi theo lô.