Sobes.tech
Junior — Middle

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.