Một quầy vé duy nhất phục vụ cả hàng người bằng cách xoay vòng thật nhanh: bán cho người này một nhịp, người kia một nhịp, ai cũng tưởng mình được phục vụ song song. Giao cho nó một khách mua một trăm vé trong một lượt — cả hàng đứng hình. Event loop của Vert.x là cái quầy đó, và client Kafka bọc client Java lại thành API trả Future một cách tiện đến mức giấu mất một chi tiết chết người: handler tiêu thụ của bạn chạy ngay trên cái quầy ấy. Nếu nó làm bất cứ việc gì tốn thời gian, thứ đứng hình không chỉ là consumer.
Kafka 3.9.0 chạy KRaft trong Docker, topic 4 phân vùng, bản tin 200 byte.
Gửi: hai cách, chênh nhau 129 lần
KafkaProducer<String,String> sx = KafkaProducer.create(vertx, cauHinh);
// cach 1: ban het roi cho
CountDownLatch l = new CountDownLatch(n);
for (int i = 0; i < n; i++)
sx.send(KafkaProducerRecord.create(topic, "k"+i, than)).onComplete(a -> l.countDown());
l.await();
// cach 2: cho xac nhan tung ban tin roi moi gui tiep
for (int i = 0; i < n; i++) {
CountDownLatch m = new CountDownLatch(1);
sx.send(...).onComplete(a -> m.countDown());
m.await();
}
| Thông lượng | Mỗi bản tin | |
|---|---|---|
| Bắn hết rồi chờ | 593 976 bản tin/giây | 1,7 µs |
| Chờ xác nhận từng bản tin | 4 613 bản tin/giây | 216,8 µs |
Chênh 129 lần, và cả hai đoạn mã đều "bất đồng bộ" theo nghĩa dùng Future. Khác biệt nằm ở chỗ cách thứ hai chờ trước khi gửi tiếp, làm producer không gom bản tin thành lô được. 216,8 µs chính là một vòng đi về tới broker; cách thứ nhất trả 1,7 µs vì hàng trăm bản tin cùng đi trong một lô.
Điều này quan trọng vì cách thứ hai rất dễ viết ra mà không nhận ra. Bất cứ khi nào bạn await, join, hay xâu chuỗi compose theo kiểu "gửi xong bản tin này mới xử lý bản tin sau", bạn đã tự đặt mình vào 4 613 bản tin/giây. Nếu thứ tự không quan trọng, đừng chờ.
Tiêu thụ: chỗ event loop bị chiếm
Đây là phần đáng đọc nhất. Tôi dựng một Verticle có cả HTTP server và Kafka consumer trên cùng một event loop — chuyện hoàn toàn bình thường trong một dịch vụ vừa nhận API vừa xử lý sự kiện. Rồi đo đường HTTP:
| Thông lượng HTTP | p50 | |
|---|---|---|
| Không có consumer | 105 425 req/s | 0,8 ms |
| Consumer xử lý nhẹ | 109 555 req/s | 0,8 ms |
| Consumer xử lý tốn CPU ngay trên event loop | 281 req/s | 377,1 ms |
| Consumer đẩy việc sang worker | 110 405 req/s | 0,8 ms |
Giảm 375 lần. Và chú ý điều gì không xảy ra: không có lỗi, không có ngoại lệ, consumer vẫn tiêu thụ bình thường. Chỉ có API HTTP của bạn đột nhiên mất 377 ms mỗi request — trong khi thủ phạm nằm ở một đoạn mã hoàn toàn khác, không liên quan gì tới HTTP, và không ai nghĩ tới khi đi tìm.
Chính consumer cũng chậm theo: cùng phép xử lý đó chạy trên event loop cho 12 546 bản tin/giây, so với 1 065 514 bản tin/giây khi handler chỉ đọc độ dài chuỗi.
Cách vá là đẩy việc ra khỏi event loop, và dừng nhận thêm trong lúc đang làm:
WorkerExecutor tho = vertx.createSharedWorkerExecutor("kafka-xu-ly", 4);
tt.handler(rec -> {
tt.pause(); // ngung nhan them
tho.executeBlocking(() -> xuLyNang(rec), false)
.onComplete(ar -> tt.resume()); // xong roi moi nhan tiep
});
Sau khi vá, HTTP quay lại 110 405 req/s — bằng đúng lúc không có consumer.
Cặp pause() / resume() là phần dễ bỏ sót. Nếu chỉ đẩy sang worker mà không pause, Vert.x vẫn tiếp tục nhận bản tin và xếp chúng vào hàng đợi của worker pool. Hàng đợi đó không có giới hạn đo được từ bên ngoài, nên khi nguồn nhanh hơn khả năng xử lý, bạn không thấy độ trễ tăng dần — bạn thấy bộ nhớ tăng dần cho tới lúc hết. Đây đúng là cùng một bài toán với vách ngăn ở phần 31: tài nguyên hữu hạn phải có ai đó giữ cửa.
Còn một hệ quả nữa của việc xử lý chậm mà số đo của tôi mới chỉ chạm tới: consumer bị coi là chết nếu quá max.poll.interval.ms (mặc định 5 phút) mà không quay lại poll. Trong lần đo của tôi đã có 2 lần phân vùng bị thu hồi. Khi điều đó xảy ra, nhóm cân bằng lại và những bản tin chưa commit sẽ được giao cho người khác — tức là xử lý lại. Nếu việc xử lý của bạn không lặp lại được an toàn, một handler chậm không chỉ làm chậm hệ thống mà còn sinh ra dữ liệu sai. Đó là lý do pause() quan trọng hơn vẻ ngoài của nó: nó giữ cho vòng poll tiếp tục chạy.
Hai chỗ vấp khi dựng Kafka để đo
Không liên quan tới Vert.x, nhưng cả hai đều tốn thời gian nên đáng ghi lại.
advertised.listeners không nhận 0.0.0.0. Khai KAFKA_LISTENERS với 0.0.0.0 làm ảnh apache/kafka suy ra advertised.listeners cũng như vậy và broker chết ngay lúc khởi động:
IllegalArgumentException: requirement failed: advertised.listeners cannot use
the nonroutable meta-address 0.0.0.0. Use a routable IP address.
Viết PLAINTEXT://:9092 thay vì PLAINTEXT://0.0.0.0:9092 là xong.
Lệnh quản trị chạy trong container lại đi vòng ra ngoài. kafka-topics.sh gọi từ bên trong container vẫn dùng địa chỉ trong advertised.listeners — ở đây là localhost:29092, cổng của máy chủ — nên nó không kết nối được:
WARN Connection to node 1 (localhost/127.0.0.1:29092) could not be established.
Broker hoàn toàn khoẻ; chỉ là địa chỉ quảng cáo được thiết kế cho client bên ngoài. Tôi bỏ hẳn lệnh shell và tạo topic bằng Admin API từ chính chương trình đo — nơi địa chỉ đó đúng.
Danh sách kiểm cho một consumer Vert.x
- Handler có làm gì quá vài chục micro giây không? Nếu có, đẩy sang worker.
- Có
pause()trước vàresume()sau không? Thiếu là đổi tràn bộ nhớ lấy tràn hàng đợi. - Có thứ gì khác chạy chung event loop không? HTTP server, timer, bất cứ Verticle nào cùng context — tất cả chết chung.
- Việc xử lý có lặp lại được an toàn không? Cân bằng lại nhóm sẽ giao lại bản tin chưa commit.
- Có đang chờ xác nhận từng bản tin lúc gửi không? Đó là 4 613 thay vì 593 976.
Muốn biết consumer có đang chiếm event loop không, không cần công cụ đo — tự bấm giờ ngay trong handler và kêu lên khi nó giữ quá lâu:
tt.handler(rec -> {
long t0 = System.nanoTime();
xuLy(rec);
long ms = (System.nanoTime() - t0) / 1_000_000;
if (ms > 10) System.out.println("handler giu event loop " + ms + " ms");
});
Dòng cảnh báo đó xuất hiện đều đặn thì mọi thứ khác trên cùng event loop đang trả giá — và vertx.executeBlocking cộng pause()/resume() là chỗ để bắt đầu.
Mẫu số chung
Cái bẫy gốc: một lớp bọc tiện tay giấu mất callback của bạn chạy trên luồng nào. API trả Future khiến tt.handler(...) trông như "cứ bất đồng bộ là an toàn", trong khi nó chạy thẳng trên event loop dùng chung — và một việc nặng ở đó đóng băng mọi thứ khác mà không một dòng lỗi nào. Câu hỏi bạn phải luôn trả lời được với mọi callback là: cái này chạy trên luồng nào, và nó có chia luồng đó với ai không? Vì thư viện hiếm khi nói ra, và mặc định gần như luôn là "trên luồng của người gọi / luồng dùng chung" — đúng cái đắt nhất. Cùng câu hỏi ấy ở một .then của Promise chạy trên main thread trình duyệt, một callback vô tình chạy trên UI thread làm khựng giao diện, một signal handler chạy trong ngữ cảnh cấm, một trigger CSDL chạy trong giao dịch của người gọi. Cái nặng thì đẩy sang chỗ riêng (worker), nhưng trước hết phải biết mình đang đứng ở đâu.
Điều thứ hai, sắc và ít ai lường: với một cơ chế "còn sống" dựa trên nhịp tim, chậm không phân biệt được với chết. Handler xử lý lâu quá max.poll.interval.ms thì Kafka tuyên bố consumer đã chết, thu hồi phân vùng, giao bản tin chưa commit cho người khác — và giờ cùng một bản tin bị xử lý hai lần. Quá tải không chỉ làm chậm; nó vượt ngưỡng nhịp tim rồi kích hoạt một cú hồi phục mà bản thân cú hồi phục ấy nhân đôi công việc và có khi kéo cả nhóm vào vòng xoáy. Cùng hình dạng ở một liveness probe của Kubernetes giết một pod chỉ đang kẹt GC, một lease/khoá phân tán hết hạn dưới tải khiến hai worker cùng chạy, một health check timeout đá một node vốn chỉ bận. Bài học: khi đặt một ngưỡng "còn sống", hãy hỏi nó có phân biệt được 'đang chậm' với 'đã chết' không — nếu không, đúng lúc hệ quá tải là lúc nó tự tay biến một cơn chậm thành một cơn chết, rồi bồi thêm việc lặp lại.
Phần sau bàn về xử lý luồng dữ liệu và áp lực ngược trong Vert.x.