Hogyan tárolja a Kafka fogyasztója az adatfolyamban elfoglalt helyét?
sobes.tech MI
Válasz az MI-től
Az Apache Kafka-ban a fogyasztó a pozícióját a adatfolyamban a tolások (offsets) segítségével tárolja. Az offset a üzenet sorszáma a téma partíciójában.
A fogyasztó időszakosan elkötelezi (commitálja) a jelenlegi offsetet, hogy újraindítás vagy hiba esetén a megfelelő helyről folytathassa az olvasást. A commit lehet automatikus (auto-commit) vagy kézi.
Go-kliens Kafka (például sarama) esetén az offset általában egy külön __consumer_offsets nevű témában van tárolva, amit a Kafka kezel. Ez lehetővé teszi, hogy több fogyasztó egy csoportban koordinálja az olvasást, és garantálja, hogy minden üzenet pontosan egyszer legyen feldolgozva.
Kézi offset commit példája Go-ban sarama használatával:
partitionConsumer, _ := consumer.ConsumePartition(topic, partition, sarama.OffsetNewest)
for msg := range partitionConsumer.Messages() {
// üzenet feldolgozása
fmt.Println(string(msg.Value))
// offset elkötelezése
consumer.MarkOffset(msg, "")
}
Így a fogyasztó pozíciója a folyamban a legutóbb elkötelezett offset, amit a Kafka tárol, és lehetővé teszi az olvasás folytatását a megfelelő helyről.