Offset фиксируется до обработки сообщения
В консьюмере сделали так, чтобы «не читать одно и то же дважды»:
for msg := range claim.Messages() {
session.MarkMessage(msg, "") // фиксация offset
if err := handler.Process(ctx, msg); err != nil {
log.Error(err)
}
}
Под был вытеснен планировщиком Kubernetes в середине батча. Что произойдёт с сообщениями, которые уже отмечены, но не обработаны?
- Они обработаются повторно после ребаланса: Kafka переотправляет всё, на что не пришла квитанция об обработке
- Ничего не изменится: отметка уходит в Kafka только при штатном завершении сессии, а при вытеснении пода пропадает
- Возникнут дубли: после ребаланса другой консьюмер группы получит партицию и перечитает тот же диапазон заново
- Они будут потеряны: группа продолжит с зафиксированного offset, и к ним уже никто не вернётся
