Cum își păstrează consumatorul în Kafka poziția în fluxul de date?
sobes.tech AI
Răspuns de la AI
În Apache Kafka, consumatorul își salvează poziția în fluxul de date folosind decalaje (offsets). Offset-ul este numărul de ordine al mesajului în partiția unui topic.
Consumatorul face periodic commit (salvare) offset-ul curent, pentru a putea relua citirea de la locul potrivit în caz de restart sau eroare. Commit-ul poate fi automat (auto-commit) sau manual.
În clienții Kafka pentru Go (de exemplu, sarama), offset-ul este de obicei salvat într-un topic special __consumer_offsets, gestionat de Kafka însuși. Acest lucru permite mai multor consumatori dintr-un grup să coordoneze citirea și să garanteze că fiecare mesaj este procesat o singură dată.
Exemplu de commit manual al offset-ului în Go folosind sarama:
partitionConsumer, _ := consumer.ConsumePartition(topic, partition, sarama.OffsetNewest)
for msg := range partitionConsumer.Messages() {
// procesarea mesajului
fmt.Println(string(msg.Value))
// commit offset
consumer.MarkOffset(msg, "")
}
Astfel, poziția consumatorului în flux este ultimul offset committat, care este stocat în Kafka și permite reluarea citirii de la locul potrivit.