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.