Sobes.tech
Junior — Middle

Kuidas Kafka tarbija salvestab oma positsiooni andmevoos?

sobes.tech AI

Vastus AI-lt

Apache Kafka-sal säilitab tarbija oma positsiooni andmevoos, kasutades nihkeid (offsets). Offset on sõnumi järjekorranumber teema partitionis.

Tarbija teeb perioodiliselt (commit) praeguse offseti, et ta saaks taaskäivitamisel või tõrke korral jätkata lugemist õige kohast. Commit võib olla automaatne (auto-commit) või käsitsi.

Go klientides Kafka (näiteks sarama) salvestatakse offset tavaliselt eraldi teemas __consumer_offsets, mida haldab Kafka ise. See võimaldab mitmel tarbijal grupis koordineerida lugemist ja tagada, et iga sõnum töödeldakse täpselt üks kord.

Näide käsitsi offseti commit'ist Go-s sarama kasutades:

partitionConsumer, _ := consumer.ConsumePartition(topic, partition, sarama.OffsetNewest)

for msg := range partitionConsumer.Messages() {
    // sõnumi töötlemine
    fmt.Println(string(msg.Value))
    // offseti commit
    consumer.MarkOffset(msg, "")
}

Seega on tarbija positsioon voolus viimati kinnitatud offset, mis on salvestatud Kafka-sse ja võimaldab lugemist jätkata õige kohast.