Sobes.tech
Junior — Middle

Как запазва потребителят в 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 и позволява да се възобнови четенето от правилното място.