Как запазва потребителят в Kafka позицията си в потока от данни?
sobes.tech AI
Отговор от AI
В Apache Kafka, потребителят съхранява позицията си в потока от данни с помощта на отмествания (offsets). Offset е поредният номер на съобщението в партицията на темата.
Потребителят периодично потвърждава (commit) текущия си offset, за да може при рестарт или срив да продължи четенето от правилното място. Commit може да бъде автоматичен (auto-commit) или ръчен.
В клиентите на Kafka за Go (например sarama), offset обикновено се съхранява в специална тема __consumer_offsets, управлявана от самия Kafka. Това позволява на няколко потребителя в група да координират четенето и да гарантират, че всяко съобщение ще бъде обработено точно веднъж.
Пример за ръчно commit-ване на offset в Go с използване на sarama:
partitionConsumer, _ := consumer.ConsumePartition(topic, partition, sarama.OffsetNewest)
for msg := range partitionConsumer.Messages() {
// обработка на съобщението
fmt.Println(string(msg.Value))
// commit на offset
consumer.MarkOffset(msg, "")
}
По този начин, позицията на потребителя в потока е последният потвърден offset, който се съхранява в Kafka и позволява да се възобнови четенето от правилното място.