Sobes.tech
Junior — Middle

Kaip Kafka vartotojas išsaugo savo poziciją duomenų sraute?

sobes.tech AI

Atsakymas iš AI

Apache Kafka, vartotojas saugo savo poziciją duomenų sraute naudodamas poslinkius (offsets). Offset yra pranešimo eilės numeris temos skiltyje.

Vartotojas periodiškai įrašo (commit) esamą offset'ą, kad perkrovus ar susidūrus su klaida, galėtų tęsti skaitymą nuo tinkamos vietos. Commit gali būti automatinis (auto-commit) arba rankinis.

Go klientuose Kafka (pvz., sarama), offset dažniausiai saugomas specialiame temoje __consumer_offsets, kurį valdo pats Kafka. Tai leidžia keliems vartotojams grupėje koordinuoti skaitymą ir užtikrinti, kad kiekvienas pranešimas būtų apdorotas tik vieną kartą.

Pavyzdys rankinio offset'ų įrašymo Go kalba naudojant sarama:

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

for msg := range partitionConsumer.Messages() {
    // pranešimo apdorojimas
    fmt.Println(string(msg.Value))
    // offset'ų įrašymas
    consumer.MarkOffset(msg, "")
}

Taip, vartotojo pozicija sraute yra paskutinis patvirtintas offset'as, kuris saugomas Kafka ir leidžia tęsti skaitymą nuo tinkamos vietos.