Chốt offset là chỗ dữ liệu thật sự mất hoặc trùng. Bài này dựng một kịch bản sập giống hệt nhau rồi chạy ba cách chốt, và đếm chính xác cái giá của từng cách.

Ba cách chốt offset, cùng một lần sập, ba kết quả khác nhau

Kịch bản

Topic một partition, 5.000 tin. Consumer xử lý mỗi tin mất 2 ms, max.poll.records=100. Chạy được 1.150 tin thì gọi Runtime.getRuntime().halt(1) — chết đột ngột, không chạy shutdown hook, không kịp chốt gì.

Sau đó bật một consumer khác cùng nhóm và đếm nó đọc được bao nhiêu.

Con số 1.150 cố ý rơi vào giữa một lô (lô đó là tin 1.101–1.200), vì sập đúng ranh giới lô lại cho kết quả khác — tôi vấp đúng chuyện này ở lần đo đầu và sẽ nói ở cuối bài.

1. Chốt tự động

p.put("enable.auto.commit", "true");
p.put("auto.commit.interval.ms", "5000");
lần 1 xử lý  1.150 tin
lần 2 xử lý  5.000 tin
tổng 6.150 cho 5.000 tin  ->  1.150 tin trùng, 0 tin mất

Consumer thứ hai đọc lại từ đầu. Không có gì được chốt cả.

Lý do: 1.150 tin × 2 ms ≈ 2,3 giây. Nhịp chốt đầu tiên ở giây thứ 5 chưa tới. Toàn bộ công đã làm bị lặp lại.

Đây là điểm quan trọng nhất về chốt tự động: cái mất không phải "vài tin cuối", mà là mọi thứ kể từ lần chốt gần nhất. Với auto.commit.interval.ms=5000 và consumer nhanh, đó có thể là hàng chục nghìn tin.

Và chốt tự động còn một tính chất ít ai để ý: nó chốt ở đầu poll(), chốt vị trí của lần poll trước. Nghĩa là nó chốt những tin bạn đã lấy về, chứ không phải những tin bạn đã xử lý xong. Nếu vòng lặp của bạn xử lý bất đồng bộ — đẩy vào một thread pool rồi quay lại poll — thì chốt tự động sẽ chốt tin còn đang nằm trong hàng đợi, và bạn có cả mất tin lẫn trùng tin cùng lúc.

2. Chốt sau khi xử lý

p.put("enable.auto.commit", "false");
// ...
for (var r : rs) xuLy(r);
c.commitSync();
lần 1 xử lý  1.150 tin
lần 2 xử lý  3.900 tin
tổng 5.050 cho 5.000 tin  ->  50 tin trùng, 0 tin mất

Offset đã chốt là 1.100 — ranh giới lô cuối cùng chốt xong. Consumer thứ hai đọc từ 1.100 đến 5.000 = 3.900 tin. Tin 1.101–1.150 được xử lý hai lần.

Số tin trùng bị chặn trên bởi max.poll.records. Đặt 100 thì tệ nhất mất công làm lại 99 tin; đặt 500 (mặc định) thì 499.

Đây là "ít nhất một lần" đúng nghĩa, và là mặc định hợp lý cho hầu hết hệ thống.

3. Chốt trước khi xử lý

var rs = c.poll(Duration.ofMillis(500));
c.commitSync();              // chốt ngay khi vừa lấy về
for (var r : rs) xuLy(r);
lần 1 xử lý  1.150 tin
lần 2 xử lý  3.800 tin
tổng 4.950 cho 5.000 tin  ->  0 tin trùng, 50 TIN MẤT HẲN

Offset chốt là 1.200 (cuối lô vừa lấy về) nhưng chỉ xử lý tới 1.150. Tin 1.151 đến 1.200 biến mất.

Và đây là phần đáng sợ: không có lỗi nào, không dòng log nào. Consumer thứ hai chạy trơn tru, báo đọc xong. Chỉ khi cộng 1.150 + 3.800 = 4.950 và so với 5.000 mới thấy thiếu.

Không ai cố tình viết như thế này. Nhưng nó xuất hiện thường xuyên dưới dạng khác: chốt trong khối finally, chốt trong một luồng riêng theo lịch, hoặc — phổ biến nhất — xử lý bất đồng bộ mà vẫn chốt đồng bộ. Mọi biến thể đều cho ra cùng kết quả: mất tin trong im lặng.

Không có cách nào cho ra đúng 5.000

Ba con số: 6.150, 5.050, 4.950. Không cái nào bằng 5.000.

Chốt trước thì mất tin. Chốt sau thì trùng tin. Kafka không cho bạn cả hai, và cũng không giả vờ là có.

Việc còn lại thuộc về phía tiêu thụ: làm cho xử lý trùng không gây hại.

-- khoá tự nhiên từ chính dữ liệu, không phải từ offset
INSERT INTO don_hang (ma_don, ...) VALUES (?, ...)
ON CONFLICT (ma_don) DO NOTHING;

Hoặc mạnh hơn — lưu offset trong cùng giao dịch với dữ liệu:

BEGIN;
  INSERT INTO don_hang ... ;
  UPDATE tien_do SET offset_da_xu_ly = ? WHERE partition = ?;
COMMIT;

Rồi khi khởi động, seek() về offset_da_xu_ly thay vì tin vào offset Kafka giữ. Cách này cho đúng một lần thật sự trên đường Kafka → cơ sở dữ liệu, vì hai việc nằm trong một giao dịch của một hệ thống.

commitSync hay commitAsync

commitSync chặn tới khi broker xác nhận, và tự thử lại khi lỗi tạm thời. commitAsync không chặn, không thử lại — nếu thử lại thì một lần chốt cũ có thể ghi đè lên một lần chốt mới hơn.

Cách dùng thường gặp: commitAsync trong vòng lặp cho nhanh, commitSync một lần trong finally lúc đóng.

try {
  while (chay) {
    var rs = c.poll(Duration.ofMillis(500));
    for (var r : rs) xuLy(r);
    c.commitAsync();
  }
} finally {
  try { c.commitSync(); } finally { c.close(); }
}

Cái commitSync cuối là thứ bảo đảm lần chốt sau cùng thật sự tới nơi trước khi tiến trình kết thúc.

Lần đo đầu của tôi không cho thấy gì

Lần chạy đầu tôi đặt điểm sập ở 1.200 tin, và cả ba cách chốt đều cho tổng đúng 5.000 — không mất, không trùng. Tôi suýt kết luận "chốt trước hay chốt sau đều như nhau".

1.200 rơi đúng ranh giới lô (max.poll.records=100). Ở đúng điểm đó, vị trí đã chốt và vị trí đã xử lý trùng nhau, nên mọi cách chốt đều đúng.

Đây là dạng phép đo sai nguy hiểm nhất: nó không hỏng, nó chỉ chọn đúng một điểm mà mọi thứ đều ổn. Dời sang 1.150 là ba con số tách hẳn ra.

Bài học chung: khi đo hành vi lúc sập, đừng để điểm sập rơi vào ranh giới tự nhiên của hệ thống. Nếu kết quả của bạn quá gọn gàng, hãy nghi ngờ nó trước đã.

Thử ba mươi giây

Xem nhóm của bạn đang chốt cách nhau bao xa:

watch -n1 'kafka-consumer-groups.sh --bootstrap-server kf:9092 \
  --describe --group nhom-cua-ban | awk "{print \$3, \$4, \$5, \$6}"'

CURRENT-OFFSET nhảy từng bước lớn cách nhau 5 giây là bạn đang dùng chốt tự động — và đó chính là lượng công sẽ phải làm lại nếu tiến trình chết ngay lúc này.

Phần sau đo phía đọc: fetch.min.bytes, max.poll.records và thông lượng consumer.