¿Cómo guarda el consumidor en Kafka la información sobre su posición en el flujo de datos?
sobes.tech AI
Respuesta de la IA
En Apache Kafka, el consumidor guarda su posición en el flujo de datos mediante desplazamientos (offsets). Offset es el número de orden del mensaje en la partición del tópico.
El consumidor realiza periódicamente un commit (guardar) del offset actual, para poder reanudar la lectura desde el lugar correcto en caso de reinicio o fallo. El commit puede ser automático (auto-commit) o manual.
En los clientes de Kafka en Go (por ejemplo, sarama), el offset generalmente se guarda en un tópico especial __consumer_offsets, gestionado por Kafka mismo. Esto permite que varios consumidores en un grupo coordinen la lectura y aseguren que cada mensaje se procese exactamente una vez.
Ejemplo de commit manual de offset en Go usando sarama:
partitionConsumer, _ := consumer.ConsumePartition(topic, partition, sarama.OffsetNewest)
for msg := range partitionConsumer.Messages() {
// procesamiento del mensaje
fmt.Println(string(msg.Value))
// commit del offset
consumer.MarkOffset(msg, "")
}
De esta forma, la posición del consumidor en el flujo es el último offset comprometido, que se almacena en Kafka y permite reanudar la lectura desde el lugar correcto.