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.