Khả năng đọc lại dữ liệu cũ là thứ phân biệt Kafka với hàng đợi truyền thống. Bài này đo bốn cách tua lại, và tìm ra một cái bẫy trong cách được dùng nhiều nhất.

Bốn cách tua lại, cái bẫy của --to-datetime, và tốc độ đọc lại

Bốn cách tua lại

Topic 300.000 tin, 3 partition, nhóm đã đọc hết:

vị trí hiện tại     p0 98.663   p1 102.317   p2 99.020
Lệnh Kết quả
--to-earliest p0 0, p1 0, p2 0
--to-latest p0 98.663, p1 102.317, p2 99.020
--shift-by -1000 p0 97.663, p1 101.317, p2 98.020
--to-offset 50000 p0 50.000, p1 50.000, p2 50.000

Chú ý dòng cuối: --to-offset đặt cùng một số cho mọi partition. Điều này hiếm khi là thứ bạn muốn — offset của các partition không liên quan gì tới nhau, và đặt tất cả về 50.000 nghĩa là bạn đọc lại nhiều hay ít tuỳ partition đó dài bao nhiêu.

--shift-by thì lùi tương đối trên từng partition, hợp lý hơn nhiều khi bạn muốn "đọc lại một nghìn tin gần nhất".

--to-datetime có một cách hỏng rất im lặng

Đây là cách được dùng nhiều nhất, vì con người nghĩ theo thời gian chứ không theo offset. Nó cũng là cách nguy hiểm nhất.

Mốc trước khi dữ liệu được ghi:

--to-datetime 2020-01-01T00:00:00.000
->  p0 0   p1 0   p2 0                    đúng

Mốc sau khi dữ liệu được ghi:

--to-datetime 2026-08-30T00:00:00.000

Warn: Partition 1 from topic rp is empty. Falling back to latest known offset.
Warn: Partition 0 from topic rp is empty. Falling back to latest known offset.
->  p0 98.663   p1 102.317   p2 99.020

Hai vấn đề chồng lên nhau.

Thông báo sai. Partition không hề rỗng — nó có hơn 100.000 tin. Ý nghĩa thật là "không có tin nào ở mốc đó trở đi". Nhưng chữ "empty" khiến người đọc nghĩ topic trống và bỏ qua.

Hành vi lùi lại là latest. Bạn yêu cầu tua về một mốc, và công cụ đưa bạn tới cuối. Nghĩa là bỏ qua toàn bộ dữ liệu — ngược hoàn toàn với ý định.

Kịch bản thực tế: có sự cố, bạn muốn xử lý lại từ 8 giờ sáng, gõ nhầm ngày thành hôm sau, chạy --execute, và consumer khởi động lại báo "không có gì để đọc, lag bằng 0". Bạn kết luận đã xử lý xong, trong khi thực ra vừa nhảy qua đúng phần dữ liệu cần xử lý.

Hai thói quen bắt buộc:

Luôn chạy --dry-run trước. Nó in ra offset mới mà không thay đổi gì.

Đọc dòng Warn. Nó không phải lỗi, không làm lệnh thất bại, và rất dễ trôi qua trong đầu ra dài.

Đọc lại nhanh đến mức nào

300.000 tin, 3 partition, một consumer:
   293,9 MB/giây     1.027.397 tin/giây

Hơn một triệu tin mỗi giây. Đọc lại toàn bộ lịch sử không đắt — phần 16 đã đo lý do: tìm tới offset ở giữa nhanh bằng tìm tới đầu, vì chỉ mục thưa cho phép nhảy thẳng tới vị trí byte.

Cái đắt nằm ở phía tiêu thụ: ghi lại vào cơ sở dữ liệu, gọi API bên ngoài, tính toán lại. Kafka đọc ra một triệu tin mỗi giây, còn consumer của bạn có thể chỉ xử lý được một nghìn.

Nên khi ước lượng thời gian xử lý lại, đừng tính theo tốc độ Kafka. Tính theo tốc độ chậm nhất trong chuỗi — thường là cơ sở dữ liệu đích.

Nhóm phải dừng hẳn

--reset-offsets --execute trả về lỗi nếu nhóm còn thành viên đang chạy.

Đây là bảo vệ có chủ ý: đổi offset dưới chân một consumer đang chạy sẽ cho kết quả không đoán được — nó có thể đã ở giữa một lô, và lần chốt tiếp theo của nó sẽ ghi đè lên giá trị bạn vừa đặt.

Quy trình đúng:

# 1. dừng mọi consumer trong nhóm, đợi tới khi rời hẳn
kafka-consumer-groups.sh --bootstrap-server kf:9092 --describe --group nhom
# CONSUMER-ID phải là "-" ở mọi dòng

# 2. thử trước
kafka-consumer-groups.sh --bootstrap-server kf:9092 --group nhom \
  --topic ten-topic --reset-offsets --to-datetime 2026-08-30T08:00:00.000 --dry-run

# 3. đọc kỹ đầu ra, kể cả dòng Warn, rồi mới chạy thật
... --execute

# 4. bật lại consumer

Ba cách xử lý lại, không chỉ một

Tua offset của nhóm hiện tại là cách thô nhất. Ba lựa chọn, theo mức độ an toàn:

Nhóm mới với auto.offset.reset=earliest. Không cần dừng gì cả: bật một nhóm mới với group.id khác, nó đọc từ đầu, và nhóm cũ chạy bình thường. Đây là cách an toàn nhất và thường là cách đúng — chỉ tốn thêm băng thông đọc.

Consumer tự seek() khi khởi động. Đọc mốc bắt đầu từ cấu hình, gọi offsetsForTimes() rồi seek(). Không phụ thuộc công cụ dòng lệnh, và bạn kiểm soát được chuyện gì xảy ra khi mốc nằm ngoài dữ liệu — thay vì nhận hành vi lùi về latest.

Map<TopicPartition, OffsetAndTimestamp> r = c.offsetsForTimes(moc);
for (var e : r.entrySet()) {
    if (e.getValue() == null) throw new IllegalStateException(
        "Khong co tin nao tu moc do tren " + e.getKey());   // dừng, đừng đoán
    c.seek(e.getKey(), e.getValue().offset());
}

Dòng throw đó chính là thứ công cụ dòng lệnh không làm. Nó biến một lỗi im lặng thành một lỗi ồn ào.

Tua offset của nhóm hiện tại. Nhanh nhất khi bạn thật sự muốn nhóm sản xuất đọc lại, nhưng đòi dừng dịch vụ và không lùi lại được sau khi --execute.

Điều kiện tiên quyết: dữ liệu phải còn

Mọi cách trên đều vô nghĩa nếu dữ liệu đã bị xoá theo retention.ms. Phần 17 đã đo: mặc định là 7 ngày, và việc dọn chạy theo chu kỳ nên có thể muộn hơn.

Nếu khả năng xử lý lại là một yêu cầu thật của hệ thống — và với hệ thống sự kiện thì thường là vậy — thì retention.ms phải được đặt theo yêu cầu đó, không phải để mặc định. "Chúng tôi cần xử lý lại được 30 ngày" là một con số cấu hình, không phải một hy vọng.

Thử ba mươi giây

Kiểm xem bạn tua lại được bao xa:

# offset sớm nhất còn giữ
kafka-get-offsets.sh --bootstrap-server kf:9092 --topic ten-topic --time -2

# dấu thời gian của tin sớm nhất đó
kafka-console-consumer.sh --bootstrap-server kf:9092 --topic ten-topic \
  --partition 0 --offset earliest --max-messages 1 \
  --property print.timestamp=true --timeout-ms 8000

Dấu thời gian đó là giới hạn thật của khả năng xử lý lại. Nếu nó gần hơn nhiều so với con số bạn hứa với ai đó, retention.ms là chỗ cần sửa — trước khi có sự cố chứ không phải trong lúc có sự cố.

Phần sau đo giao dịch: gói một tin vào một giao dịch cho ra 96 tin mỗi giây.