Khi một tin không xử lý được, phản xạ tự nhiên là thử lại. Bài này đo cái giá của phản xạ đó, và nó lớn hơn nhiều so với vẻ ngoài.
Đo
3.000 tin, mỗi tin xử lý tốn 1 ms. Tin hỏng thì thử lại 3 lần, mỗi lần chờ 20 ms — một cơ chế thử lại rất khiêm tốn so với những gì thường thấy trong mã sản xuất.
| Tỉ lệ hỏng | Thử lại tại chỗ | Đẩy sang DLQ | DLQ nhanh hơn |
|---|---|---|---|
| 1% | 350 tin/giây | 448 | 1,3× |
| 5% | 177 | 469 | 2,6× |
| 20% | 72 | 538 | 7,5× |
| 50% | 30 | 766 | 25,5× |
Hai chiều ngược nhau.
Cột "thử lại" sụp theo tỉ lệ hỏng: 350 xuống 30, tức mười một lần.
Cột DLQ thì nhanh lên: 448 lên 766. Vì tin hỏng bỏ qua bước xử lý 1 ms và chỉ bị chuyển tiếp — càng nhiều tin hỏng thì càng ít việc thật phải làm.
Ở mức 50%, chênh lệch là 25,5 lần.
Vì sao thử lại tại chỗ đắt đến vậy
Một partition được xử lý bởi đúng một consumer, tuần tự. Ngồi chờ một tin hỏng nghĩa là mọi tin sau nó cũng ngồi chờ.
offset 100 hỏng -> chờ 3 × 20 ms = 60 ms
offset 101 tốt -> phải đợi 60 ms đó trước khi tới lượt
offset 102 tốt -> cũng vậy
Ở tỉ lệ 5%, cứ 20 tin lại có một tin làm 19 tin kia trễ thêm 60 ms. Đây là hàng đợi bị chặn ở đầu, và Kafka không có cách nào vượt qua trong cùng một partition — đó chính là điều nó bảo đảm về thứ tự.
Con số 20 ms trong phép đo này còn rất nhẹ. Cơ chế thử lại thật thường có độ lùi tăng dần: 1 giây, 2 giây, 4 giây. Ba lần thử với độ lùi đó là 7 giây cho mỗi tin hỏng, và ở tỉ lệ 5% thì thông lượng gần như về không.
Tệ hơn nữa: max.poll.interval.ms mặc định là 300.000 ms. Nếu tổng thời gian chờ của một lô vượt qua đó, broker coi consumer đã chết và cân bằng lại nhóm — phần 9 đã đo cái giá của việc đó. Lô bị bỏ dở được giao cho consumer khác, gặp đúng tin hỏng, chờ lại từ đầu, và vòng lặp tiếp tục.
Kiến trúc ba tầng
Cách xử lý đúng là chuyển việc chờ ra khỏi đường đi chính:
don-hang topic chính, xử lý ngay
| hỏng
v
don-hang-retry-5s consumer riêng, đợi 5 giây rồi thử lại
| vẫn hỏng
v
don-hang-retry-1m đợi 1 phút rồi thử lại
| vẫn hỏng
v
don-hang-dlq không tự xử lý nữa, chờ người xem
Mỗi tầng là một topic riêng với consumer riêng. Việc chờ ở tầng dưới không chặn topic chính, vì chúng là những vòng lặp poll hoàn toàn khác nhau.
Tin mang theo header đếm số lần đã thử:
var h = r.headers().lastHeader("so-lan-thu");
int lan = h == null ? 0 : Integer.parseInt(new String(h.value()));
var moi = new ProducerRecord<>(topicKe, r.key(), r.value());
moi.headers().add("so-lan-thu", String.valueOf(lan + 1).getBytes());
moi.headers().add("loi-goc", e.toString().getBytes());
moi.headers().add("topic-goc", r.topic().getBytes());
producer.send(moi);
Ba header đó là tối thiểu. Không có loi-goc thì khi mở DLQ ra bạn có một đống tin mà không biết vì sao chúng ở đó, và phải đoán.
Giữ nguyên khoá gốc khi chuyển tiếp. Nếu không, thứ tự theo khoá bị phá vỡ ở tầng thử lại — và với dữ liệu cần thứ tự, việc thử lại có thể tạo ra kết quả sai thay vì chỉ chậm.
Cách chờ trong tầng thử lại
Đây là chỗ dễ làm sai. Thread.sleep(5000) trong vòng lặp poll của tầng retry-5s cũng chặn chính tầng đó — bạn chỉ dời vấn đề đi một bước.
Cách đúng là tạm dừng partition rồi tiếp tục poll:
c.pause(c.assignment());
// vẫn gọi poll() để giữ nhịp tim, nhưng nó trả về rỗng
while (chuaDenGio()) c.poll(Duration.ofMillis(500));
c.resume(c.assignment());
poll() khi đang tạm dừng vẫn gửi nhịp tim và vẫn tính vào max.poll.interval.ms, nên nhóm không cân bằng lại. Đây là khác biệt then chốt giữa "chờ đúng cách" và "chờ rồi bị đá khỏi nhóm".
Cách khác, đơn giản hơn: nhìn dấu thời gian của tin, nếu chưa đủ 5 giây thì seek lùi lại và đọc lô sau. Không giữ trạng thái gì, nhưng đọc lặp nhiều hơn.
Khi nào KHÔNG nên thử lại
Thử lại chỉ có nghĩa với lỗi tạm thời: cơ sở dữ liệu bận, dịch vụ bên ngoài trả 503, mạng chập chờn.
Với lỗi vĩnh viễn thì thử lại là lãng phí thuần tuý:
- Tin không giải mã được (phần 22) — thử một triệu lần vẫn không giải mã được.
- Trường bắt buộc bị thiếu.
- Tham chiếu tới bản ghi không tồn tại.
- Lỗi nghiệp vụ: đơn hàng đã huỷ, tài khoản đã đóng.
Phân loại lỗi trước khi quyết định, và đẩy thẳng lỗi vĩnh viễn sang DLQ:
catch (LoiTamThoi e) { chuyenSang(topicRetry, r, e); }
catch (Exception e) { chuyenSang(topicDlq, r, e); } // mọi thứ còn lại
Mặc định nên là DLQ, không phải thử lại. Thử lại là ngoại lệ dành cho những lỗi bạn đã xác định rõ là tạm thời.
Hàng đợi chết cần được nhìn
Một DLQ không ai xem là một cách đánh mất dữ liệu chậm rãi. Ba thứ tối thiểu:
Cảnh báo theo tốc độ vào, không theo tổng số. Tổng số chỉ tăng và sẽ vượt mọi ngưỡng cố định. Cái đáng báo là "mười tin vào DLQ trong một phút" — nghĩa là có gì đó vừa hỏng.
Đặt retention.ms dài cho topic DLQ. Nếu DLQ hết hạn trong 7 ngày như mặc định, bằng chứng của sự cố biến mất trước khi có ai kịp điều tra.
Có cách đẩy tin quay lại. Sau khi sửa lỗi, bạn cần một công cụ đọc DLQ và ghi lại vào topic gốc. Viết nó trước khi cần dùng, không phải lúc đang có sự cố.
Thử ba mươi giây
Xem đường ống của bạn có DLQ không, và nó có ai xem không:
# có topic DLQ nào không
kafka-topics.sh --bootstrap-server kf:9092 --list | grep -iE "dlq|dead|error|retry"
# nếu có, xem nó chứa bao nhiêu và tin mới nhất từ bao giờ
kafka-get-offsets.sh --bootstrap-server kf:9092 --topic ten-dlq
kafka-console-consumer.sh --bootstrap-server kf:9092 --topic ten-dlq \
--offset -1 --partition 0 --max-messages 1 --property print.timestamp=true --timeout-ms 8000
Nếu tin mới nhất trong DLQ có từ nhiều tháng trước và không ai biết, bạn vừa tìm thấy một đống dữ liệu đã lặng lẽ rơi ra khỏi hệ thống.
Phần sau đo Kafka Connect: đưa dữ liệu vào và ra Kafka mà không viết mã.