"Ghi đúng một lần" là cụm từ dễ hiểu sai nhất trong Kafka. Bài này mở tệp log ra xem cơ chế thật, và tìm được một hệ quả vận hành ít ai nhắc tới.

Chi phí idempotence, dấu vết trên đĩa, và giao dịch bị huỷ

Vì sao cần

Producer gửi lô, mạng đứt, không nhận được phản hồi. Nó không biết broker đã ghi hay chưa. Hai lựa chọn:

  • Không thử lại — mất tin nếu broker chưa ghi.
  • Thử lại — trùng tin nếu broker đã ghi.

Trước Kafka 0.11 chỉ có hai lựa chọn đó. enable.idempotence thêm lựa chọn thứ ba.

Cơ chế, nhìn trên đĩa

Producer bật idempotence xin broker cấp một producerId, rồi đánh số thứ tự cho từng lô của từng partition. Broker nhớ số cuối cùng đã nhận cho mỗi cặp (producerId, partition). Lô gửi lại mang số đã thấy thì bị bỏ qua — và vẫn trả về thành công, nên producer không cần biết gì.

kafka-dump-log.sh cho thấy dấu vết đó ngay trong phần đầu mỗi lô:

# tắt idempotence
baseOffset: 0 lastOffset: 15 count: 16 baseSequence: -1 producerId: -1 producerEpoch: -1

# bật
baseOffset: 0 lastOffset: 4 count: 5 baseSequence: 0 lastSequence: 4 producerId: 23 producerEpoch: 0

producerEpoch là thứ xử lý trường hợp producer chết rồi sống lại: bản mới xin epoch cao hơn, và broker từ chối mọi lô mang epoch cũ. Đó là cách một producer "thây ma" bị chặn.

Hai giới hạn cần nhớ rõ, vì chúng là nguồn của phần lớn hiểu nhầm:

  1. Chỉ trong phạm vi một phiên producer. Producer khởi động lại xin producerId mới, và mọi thứ nó gửi lại đều là tin mới đối với broker.
  2. Chỉ chống trùng do chính producer thử lại. Ứng dụng của bạn gọi send() hai lần cho cùng một đơn hàng thì đó là hai tin khác nhau, và không có gì chặn.

Chi phí

300.000 tin × 500 byte, acks=all, ba lần mỗi mức:

Thông lượng Độ trễ trung bình
tắt 343.642 / 513.698 / 476.947 13,46 / 1,18 / 2,08 ms
bật 413.793 / 397.350 / 466.562 0,73 / 0,63 / 0,49 ms

Thông lượng của hai nhóm chồng lấn nhau — không tách được bằng phép đo này. Độ trễ thì nhóm bật thấp hơn ổn định qua cả ba lần.

Từ Kafka 3.0 nó bật mặc định. Với con số này, không có lý do gì để tắt.

Một điều kiện đi kèm: max.in.flight.requests.per.connection phải ≤ 5. Đặt cao hơn thì producer từ chối khởi động. Con số 5 là số lô broker giữ lại để nhận ra bản trùng; vượt qua là nó không còn theo dõi nổi thứ tự.

Giao dịch, và cái nằm lại trên đĩa

Idempotence lo một lô một partition. Giao dịch lo nhiều partition, nhiều topic, một lần cam kết — hoặc tất cả hiện ra, hoặc không cái nào.

Properties p = new Properties();
p.put("bootstrap.servers", "kf:9092");
p.put("transactional.id", "tx-demo-2");
p.put("enable.idempotence", "true");
KafkaProducer<String,String> pr = new KafkaProducer<>(p);
pr.initTransactions();

pr.beginTransaction();
for (int i = 0; i < 5; i++) pr.send(new ProducerRecord<>("tx2", "k"+i, "COMMIT-"+i));
pr.commitTransaction();

pr.beginTransaction();
for (int i = 0; i < 5; i++) pr.send(new ProducerRecord<>("tx2", "k"+i, "ABORT-"+i));
pr.flush();
pr.abortTransaction();

Năm tin cam kết, năm tin huỷ. Đọc lại:

isolation.level=read_uncommitted   ->  10 tin  (cả COMMIT lẫn ABORT)
isolation.level=read_committed     ->   5 tin  (chỉ COMMIT)

Đúng như mong đợi. Nhưng nhìn vào đĩa:

offset  0..4     COMMIT-0 .. COMMIT-4
offset  5        bản ghi điều khiển   isControl=true
offset  6..10    ABORT-0 .. ABORT-4        <- vẫn ghi đầy đủ
offset 11        bản ghi điều khiển   isControl=true

Huỷ giao dịch không xoá gì. Năm tin bị huỷ nằm nguyên trên đĩa, chiếm offset 6 đến 10. Broker chỉ ghi thêm một bản ghi điều khiển đánh dấu "giao dịch này huỷ", và consumer mới là nơi lọc chúng ra.

Hai bản ghi điều khiển ở offset 5 và 11 cũng chiếm offset thật, dù không consumer nào đọc được chúng.

Hệ quả vận hành ít ai nói

Topic tx2offset cuối là 12 nhưng chỉ có 5 tin đọc được.

Độ trễ consumer được tính bằng offset cuối − offset đã cam kết. Trên topic có giao dịch, consumer đã đọc hết mọi thứ nó đọc được vẫn cam kết ở offset 5, trong khi offset cuối là 12.

Lag hiện 7, và nó sẽ mãi là 7.

Đây là nguồn của rất nhiều cảnh báo giả. Nếu bạn đặt cảnh báo ở lag > 0 trên một topic dùng giao dịch, nó sẽ kêu suốt ngày. Và tệ hơn: khi có sự cố thật thì không ai còn để ý nữa.

Cách đúng: đặt ngưỡng theo mức lag tăng dần theo thời gian, hoặc theo records-lag-max mà consumer tự báo qua JMX, chứ không theo hiệu số offset thô.

Tỉ lệ này lớn hơn bạn nghĩ: một giao dịch chỉ chứa một tin sẽ tạo một bản ghi điều khiển đi kèm — 50% offset là điều khiển. Giao dịch càng nhỏ, con số lag cố định càng lớn.

"Ghi đúng một lần" nghĩa là gì

Đây là chỗ cụm từ này gây hiểu nhầm nhiều nhất. Kafka bảo đảm đúng một lần cho đường đi Kafka → xử lý → Kafka: đọc từ topic, tính toán, ghi sang topic khác, và cam kết offset đọc trong cùng giao dịch với việc ghi. Kafka Streams làm sẵn chuyện này khi bật processing.guarantee=exactly_once_v2.

không bảo đảm đúng một lần cho:

  • Kafka → cơ sở dữ liệu. Ghi vào PostgreSQL rồi cam kết offset là hai hệ thống, không có giao dịch chung. Cách xử lý là làm phía tiêu thụ bất biến với lặp lại: dùng khoá tự nhiên và INSERT ... ON CONFLICT DO NOTHING, hoặc lưu offset đã xử lý ngay trong cùng bảng và cùng giao dịch với dữ liệu.
  • Kafka → gọi HTTP. Không có cách nào. Phải để phía nhận nhận diện bản trùng, thường bằng khoá idempotency trong header.

Cụm từ chính xác cho phần lớn hệ thống thật là "ít nhất một lần, cộng với phía tiêu thụ chịu được lặp". Nó kém hào nhoáng hơn nhưng dựng được, còn "đúng một lần" xuyên qua ranh giới hệ thống thì không.

Chi phí của giao dịch

Giao dịch tốn hơn idempotence rõ rệt: mỗi lần cam kết là một vòng trao đổi với điều phối viên giao dịch, cộng thêm bản ghi điều khiển trên mọi partition đã chạm tới. Gộp nhiều tin vào một giao dịch thay vì mỗi tin một giao dịch là cách duy nhất giữ chi phí đó ở mức chấp nhận được.

Và consumer đọc read_committed phải đợi tới khi giao dịch kết thúc mới thấy tin — nghĩa là độ trễ đầu-cuối gắn với thời gian giao dịch của bạn, không gắn với thời gian gửi.

Thử ba mươi giây

Kiểm chuyện offset trên chính topic của bạn:

# offset cuối
kafka-get-offsets.sh --bootstrap-server kf:9092 --topic topic-cua-ban

# số tin đọc được thật
kafka-console-consumer.sh --bootstrap-server kf:9092 --topic topic-cua-ban \
  --from-beginning --timeout-ms 10000 --consumer-property isolation.level=read_committed | wc -l

Hai con số bằng nhau: topic không dùng giao dịch. Lệch nhau: phần chênh là bản ghi điều khiển và tin bị huỷ, và đó chính là mức lag mà consumer của bạn không bao giờ hạ xuống được.

Phần sau chuyển sang phía đọc: consumer group chia partition thế nào, và cân bằng lại tốn bao lâu.