Pregunta principal : ejecutamos aplicaciones Kafka Streams (Java) en Kubernetes para consumir, procesar y producir datos en tiempo real en nuestro Kafka Cluster (ejecutando Confluent Community Edition v7.0/Kafka v3.0). ¿Cómo podemos hacer una implementación de nuestras aplicaciones de una manera que limite el tiempo de inactividad en el consumo de registros? Nuestro objetivo inicial era aproximadamente 2 sec tiempo de inactividad una sola vez para cada tarea.
Nuestro objetivo es realizar implementaciones continuas de cambios en el entorno de producción, pero las implementaciones son demasiado disruptivas al causar tiempo de inactividad en el consumo de registros en nuestras aplicaciones, lo que genera latencia en los registros producidos en tiempo real.
Hemos probado diferentes estrategias para ver cómo afecta la latencia (tiempo de inactividad).
Estrategia #1:
85 secEstrategia #2:
3 minutes para permitir la restauración del estado local en la nueva instancia de la aplicación3 minutes termine una instancia de aplicación anterior39 secEstrategia #3:
15 minutes7 sec . Sin embargo 15 minutes por instancia de la aplicación generarán 15 minutos x 6 instancias = 90 minutes para implementar un cambio + 30 minutes adicionales para que finalice el protocolo de reequilibrio incremental. Consideramos que el tiempo de implementación es bastante excesivo.Hemos estado leyendo el KIP-429: Protocolo de reequilibrio incremental del consumidor de Kafka e intentamos configurar la aplicación para admitir nuestro caso de uso.
A continuación se muestra la configuración clave de Kafka Streams que hicimos para las estrategias n.º 2 y n.º 3:
acceptable.recovery.lag: 6000 num.standby.replicas: 2 max.warmup.replicas: 6 probing.rebalance.interval.ms: 60000 num.stream.threads: 2 El tema de entrada tiene 12 partitions y una tasa de mensajes de 800 records/s en promedio. Hay 3 almacenes de estado de valores clave de Kafka Streams, donde dos tienen la misma tasa como tema de entrada. Estos dos tienen un 20GB size . El tamaño del juego de llaves es de aproximadamente 4000 . En teoría, acceptable.recovery.lag anterior debería tener una latencia de ~60 segundos en el tema del registro de cambios por partición.
Estas son algunas métricas (aumento, tasa de mensajes recibidos y latencia de los mensajes recibidos) por instancia de aplicación para la estrategia n.º 3 : 
Observaciones notables que hicimos:
1a : comienza la primera instancia nueva de la aplicación
1b : Kafka se reequilibra inmediatamente, se asignan más tareas a 2 instancias de aplicaciones antiguas y se pierden 2 tareas antiguas
1c : al mismo tiempo, la latencia máxima registrada aumenta de 0.2 sec a 3.5 sec (esto indica que el reequilibrio tarda aproximadamente 3 sec )
2 : se produce un reequilibrio de sondeo y Kafka Streams de alguna manera decide revocar una tarea de una de las instancias de la aplicación anterior y dársela a una que ya tiene la mayoría de las tareas
3 : se finaliza la última instancia de la aplicación anterior
4 : todas las particiones se reequilibran como antes de la actualización, el reequilibrio incremental se completa (aproximadamente 33 minutes después de la finalización de la última instancia de la aplicación)
Other : se tarda aproximadamente 40 minutes en asignar una tarea a la primera instancia nueva de la aplicación. Además, cada tarea se reasigna varias veces, lo que provoca muchas pequeñas interrupciones de 3 sec .
Podemos proporcionar muchos más detalles sobre la topología, los temas, la configuración y los diagramas de métricas para las otras estrategias, si es necesario (este hilo ya es enorme).
Algunos consejos generales para su aplicación Kafka Stream (no es un profesional, pero es algo que observo personalmente).
group.instance.id: "${hostname}" . De esta manera, el pod mantendrá el mismo nombre de pod y podrá usar el Protocolo de reequilibrio incremental del consumidor de Kafka.terminationGracePeriodSeconds para permitir el apagado adecuado. De lo contrario, Kafka Stream detectará que el punto de control no está limpio y recuperará el almacén de estados (no está muy claro si su almacén de estados persiste, pero creo que ya está hecho).Lo único que he observado por mi parte es que un simple reequilibrio del grupo de consumidores es más largo que los 2 segundos esperados y parece ser un objetivo complicado (como dijiste, un reequilibrio toma 3 segundos). Pero al usar el reequilibrio incremental, creo que tendrá alguna partición que se procesará durante el reequilibrio (no lo he medido personalmente por ahora).