Comment le consommateur dans Kafka enregistre-t-il sa position dans le flux de données?
sobes.tech IA
Réponse de l'IA
Dans Apache Kafka, le consommateur enregistre sa position dans le flux de données à l’aide de décalages (offsets). L’offset est le numéro de l’ordre du message dans la partition du topic.
Le consommateur effectue périodiquement un commit (enregistrement) de l’offset actuel, afin de pouvoir reprendre la lecture à partir du bon endroit en cas de redémarrage ou de panne. Le commit peut être automatique (auto-commit) ou manuel.
Dans les clients Kafka en Go (par exemple, sarama), l’offset est généralement stocké dans un topic spécial __consumer_offsets, géré par Kafka lui-même. Cela permet à plusieurs consommateurs d’un groupe de coordonner la lecture et de garantir que chaque message est traité une seule fois.
Exemple de commit manuel d’offset en Go avec sarama:
partitionConsumer, _ := consumer.ConsumePartition(topic, partition, sarama.OffsetNewest)
for msg := range partitionConsumer.Messages() {
// traitement du message
fmt.Println(string(msg.Value))
// commit de l’offset
consumer.MarkOffset(msg, "")
}
Ainsi, la position du consommateur dans le flux est le dernier offset committé, stocké dans Kafka, permettant de reprendre la lecture à partir du bon endroit.