Sobes.tech
Junior — Middle

Jak konsument w Kafka zapisuje swoją pozycję w strumieniu danych?

sobes.tech AI

Odpowiedź od AI

W Apache Kafka konsument zapisuje swoją pozycję w strumieniu danych za pomocą przesunięć (offsetów). Offset to numer porządkowy wiadomości w partycji tematu.

Konsument okresowo zatwierdza (commituje) aktualny offset, aby w przypadku ponownego uruchomienia lub awarii móc kontynuować czytanie od odpowiedniego miejsca. Commit może być automatyczny (auto-commit) lub ręczny.

W klientach Kafka dla Go (np. sarama) offset jest zwykle przechowywany w specjalnym temacie __consumer_offsets, zarządzanym przez Kafka. Umożliwia to kilku konsumentom w grupie koordynację odczytu i gwarantuje, że każda wiadomość zostanie przetworzona dokładnie raz.

Przykład ręcznego zatwierdzania offsetu w Go z użyciem sarama:

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

for msg := range partitionConsumer.Messages() {
    // przetwarzanie wiadomości
    fmt.Println(string(msg.Value))
    // zatwierdzenie offsetu
    consumer.MarkOffset(msg, "")
}

W ten sposób pozycja konsumenta w strumieniu to ostatni zatwierdzony offset, który jest przechowywany w Kafka i pozwala wznowić odczyt z odpowiedniego miejsca.