Sobes.tech
Junior — Middle

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.