Tengo una aplicación Spring Boot en la que usamos Spring Kafka en la sección del consumidor. He hecho enable.auto.commit en falso y configuré el listener ack-mode to manual_immediate
Tengo consumidores simultáneos, así que después de consumir el registro, llamo acknowledgment.acknowledge() pero aquí todavía me enfrento al problema del problema duplicado cada vez que ocurre un reequilibrio, otro consumidor comienza a consumir el mismo mensaje que ya ha consumido un consumidor. Alguna idea de qué magia está sucediendo detrás de escena.
¿Alguien sabe cuando se usa manual_immediate, se confirma el mensaje mediante commitSync o commitAsync? ¿Hay alguna forma en que podamos cambiar el comportamiento para evitar la lectura duplicada de mensajes de registro? ¿Hay alguna manera de que podamos usar el modelo híbrido en Spring Kafka?
En Spring Boot Kafka, hay una manera en que podemos ver cada vez que ocurre un reequilibrio, puedo registrarlo.
¿Cómo crear un reequilibrio si queremos hacerlo con algún propósito de prueba?
Siempre que llame al acknowledge en el hilo del oyente, utilizará commitSync() de forma predeterminada; use la propiedad de contenedor syncCommits para usar confirmaciones asíncronas.
Si lo llama en un subproceso diferente, la confirmación se pone en cola para que el subproceso del consumidor la procese lo antes posible.
No se pueden evitar los duplicados si se produce un reequilibrio forzado porque su oyente tardó demasiado en procesar los registros recibidos por poll() .
Puede aumentar max.poll.interval.ms y/o reducir max.poll.records para asegurarse de que puede procesar los registros a tiempo.
Puede agregar un ConsumerRebalanceListener a las propiedades del contenedor para registrar reequilibrios.
Reduzca max.poll.interval.ms a un valor pequeño para reproducirlo en una prueba.
En primer lugar, independientemente del modo de acuse de recibo, nunca se garantiza que un mensaje se consuma solo una vez. Por ejemplo, puede ocurrir un reequilibrio entre el momento en que se consume un mensaje y se confirma la compensación de tiempo, lo que hace que Kafka entregue el mensaje nuevamente al consumidor recién asignado. Es responsabilidad de las aplicaciones ser idempotentes ante los mensajes duplicados.
Para escuchar eventos de reequilibrio, se necesita una implementación de ConsumerRebalanceListener . Puede conectar esta implementación a la instancia ConcurrentKafkaListenerContainerFactory configurada automáticamente de Spring. Una descripción más detallada de cómo se puede hacer esto ya ha sido respondida aquí .
Si desea crear un reequilibrio forzado para la prueba, puede hacerlo matando a uno de los más de 1 consumidores existentes. Si usa spring-kafka, puede hacerlo usando una instancia @Autowired de KafkaListenerEndpointRegistry y matar/descansar/(re)iniciar cualquier consumidor. Algo sobre esto debería hacer:
@Autowired KafkaListenerEndpointRegistry registry; public void myTest() { Collection<MessageListenerContainer> containers = registry.getAllListenerContainers() containers.get(0).stop() }