Tengo una aplicación de Kafka Streams que consume y produce en un clúster de Kafka con 3 intermediarios y un factor de replicación de 3. Aparte de los temas de compensación del consumidor (50 particiones), todos los demás temas tienen solo una partición cada uno.
Cuando los intermediarios intentan elegir una réplica preferida, la aplicación Streams (que se ejecuta en una instancia completamente diferente a la de los intermediarios) falla con el error:
Caused by: org.apache.kafka.streams.errors.StreamsException: task [0_0] exception caught when producing at org.apache.kafka.streams.processor.internals.RecordCollectorImpl.checkForException(RecordCollectorImpl.java:119) ... at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:197) Caused by: org.apache.kafka.common.errors.NotLeaderForPartitionException: This server is not the leader for that topic-partition.¿Es normal que la aplicación Streams intente ser la líder de la partición, dado que se ejecuta en un servidor que no forma parte del clúster de Kafka?
Puedo reproducir este comportamiento a pedido por:
bin/kafka-preferred-replica-election.sh --zookeeper localhostMi problema parece ser similar a este error informado , por lo que me pregunto si se trata de un nuevo error de Kafka Streams. Mi seguimiento completo de la pila es literalmente exactamente el mismo que el vínculo esencial en el error informado ( here ).
Otro detalle potencialmente interesante es que durante la elección del líder, recibo estos mensajes en el controller.log del corredor:
[2017-04-12 11:07:50,940] WARN [Controller-3-to-broker-3-send-thread], Controller 3's connection to broker BROKER-3-HOSTNAME:9092 (id: 3 rack: null) was unsuccessful (kafka.controller.RequestSendThread) java.io.IOException: Connection to BROKER-3-HOSTNAME:9092 (id: 3 rack: null) failed at kafka.utils.NetworkClientBlockingOps$.awaitReady$1(NetworkClientBlockingOps.scala:84) at kafka.utils.NetworkClientBlockingOps$.blockingReady$extension(NetworkClientBlockingOps.scala:94) at kafka.controller.RequestSendThread.brokerReady(ControllerChannelManager.scala:232) at kafka.controller.RequestSendThread.liftedTree1$1(ControllerChannelManager.scala:185) at kafka.controller.RequestSendThread.doWork(ControllerChannelManager.scala:184) at kafka.utils.ShutdownableThread.run(ShutdownableThread.scala:63)Inicialmente pensé que este error de conexión era el culpable, pero después de que la elección del líder bloquea la aplicación Streams, si reinicio la aplicación Streams, funciona normalmente hasta la próxima elección, sin que yo toque a los corredores en absoluto.
Todos los servidores (3 agentes Kafka y la aplicación Streams) se ejecutan en instancias EC2.
Esto ahora está arreglado en 0.10.2.1. Si no puede retomar eso, asegúrese de tener estos dos parámetros configurados de la siguiente manera en la configuración de sus transmisiones:
final Properties props = new Properties(); ... props.put(ProducerConfig.RETRIES_CONFIG, 10); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, Integer.toString(Integer.MAX_VALUE));