Client Kafka của Vert.x bọc client Java chính thức lại thành API trả Future. Cái bọc đó tiện, và nó cũng giấu đi một chi tiết quan trọng: handler tiêu thụ của bạn chạy ngay trên event loop. Nếu chỗ đó làm bất cứ việc gì tốn thời gian, thứ chết 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 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.
  • 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.

Thử ba mươi giây

Kiểm tra xem consumer của bạn có đang chiếm event loop không, không cần công cụ đo:

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.

Phần sau bàn về xử lý luồng dữ liệu và áp lực ngược trong Vert.x.