El problema:
Tengo tres instancias de una aplicación Java ejecutándose en Kubernetes. Mi aplicación usa Apache Camel para leer desde un flujo de Kinesis. Actualmente estoy observando dos problemas relacionados:
Cada una de las tres instancias en ejecución de mi aplicación está procesando los registros que ingresan a la secuencia, cuando solo quiero que cada registro se procese una vez (quiero tres en funcionamiento para fines de escala). Tenía la esperanza de que mientras una instancia está procesando el registro A, una segunda podría estar recogiendo el registro B, etc.
Cada vez que mi aplicación se vuelve a implementar en Kubernetes, cada instancia inicia cada registro nuevamente (en otras palabras, no tiene idea de dónde se quedó o qué registros se procesaron previamente).
Después de 5 minutos, el iterador de fragmentos que usa mi aplicación para sondear Kinesis se agota. Sé que este es un comportamiento normal, pero lo que no entiendo es por qué mi aplicación no obtiene un nuevo iterador. Esta captura de pantalla muestra el error de DataDog. 
Lo que probé: en primer lugar, creo que este problema se debe a identificadores de iteradores de fragmentos e identificadores de consumidores de Kinesis inconsistentes en las tres instancias de mi aplicación y en todas las implementaciones. Sin embargo, no he podido ubicar dónde se establecen estos valores en el código y cómo podría configurarlos. Por supuesto, también puede haber una solución mejor. He encontrado muy poca documentación sobre Kinesis/Kubernetes/Camel trabajando juntos, y muy pocas fuentes externas han sido útiles.
La documentación sobre AWS Kinesis :: Apache Camel es muy limitada, pero he intentado jugar con el tipo de iterador y crear una configuración de cliente personalizada.
Avíseme si necesita información adicional, gracias.
Configurando el cliente:
main.bind("kinesisClient", AmazonKinesisClientBuilder.defaultClient()); . . . inputUri = String.format("aws-kinesis://%s?amazonKinesisClient=#kinesisClient", rawKinesisName); main.configure().addRoutesBuilder(new RawDataRoute(inputUri, inputTransform));Mi ruta:
public class RawDataRoute extends RouteBuilder { private static final Logger LOG = new Logger(RawDataRoute.class, true); private String rawDataStreamUri; private Expression transform; public RawDataRoute(final String rawDataStreamUri, final Expression transform) { this.rawDataStreamUri = rawDataStreamUri; this.transform = transform; } @Override public void configure() { // TODO add error handling from(rawDataStreamUri) .routeId("raw_data_stream") .transform(transform) .to("direct:main_input_stream"); } }