Tengo el siguiente código
class Consumer(val consumer: KafkaConsumer<String, ConsumerRecord<String>>) { fun run() { consumer.seekToEnd(emptyList()) val pollDuration = 30 // seconds while (true) { val records = consumer.poll(Duration.ofSeconds(pollDuration)) // perform record analysis and commitSync() } } } } El tema al que está suscrito el consumidor recibe continuamente registros. Ocasionalmente, el consumidor se bloqueará debido al paso de procesamiento. Cuando el consumidor se reinicia, quiero que consuma desde el último desplazamiento en el tema (es decir, ignore los registros que se publicaron en el tema mientras el consumidor estaba inactivo). Pensé que el método seekToEnd() lo aseguraría. Sin embargo, parece que el método no tiene ningún efecto en absoluto. El consumidor comienza a consumir desde el desplazamiento desde el que se estrelló.
¿Cuál es la forma correcta de usar seekToEnd() ?
Editar: el consumidor se crea con las siguientes configuraciones
fun <T> buildConsumer(valueDeserializer: String): KafkaConsumer<String, T> { val props = setupConfig(valueDeserializer) Common.setupConsumerSecurityProtocol(props) return createConsumer(props) } fun setupConfig(valueDeserializer: String): Properties { // Configuration setup val props = Properties() props[ConsumerConfig.GROUP_ID_CONFIG] = config.applicationId props[ConsumerConfig.CLIENT_ID_CONFIG] = config.kafka.clientId props[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = config.kafka.bootstrapServers props[AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG] = config.kafka.schemaRegistryUrl props[ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG] = config.kafka.stringDeserializer props[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = valueDeserializer props[KafkaAvroDeserializerConfig.SPECIFIC_AVRO_READER_CONFIG] = "true" props[ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG] = config.kafka.maxPollIntervalMs props[ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG] = config.kafka.sessionTimeoutMs props[ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG] = "false" props[ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG] = "false" props[ConsumerConfig.AUTO_OFFSET_RESET_CONFIG] = "latest" return props } fun <T> createConsumer(props: Properties): KafkaConsumer<String, T> { val consumer = KafkaConsumer<String, T>(props) consumer.subscribe(listOf(config.kafka.inputTopic)) return consumer }¡Encontré una solución!
Necesitaba agregar una encuesta ficticia como parte del proceso de inicialización del consumidor. Dado que varios métodos de Kafka se evalúan con pereza, es necesario realizar una encuesta ficticia para asignar particiones al consumidor. Sin la encuesta ficticia, el consumidor intenta buscar hasta el final de las particiones que son nulas. Como resultado, seekToEnd() no tiene efecto.
Es importante que la duración de la encuesta ficticia sea lo suficientemente larga para que se asignen las particiones. Por ejemplo, con consumer.poll((Duration.ofSeconds(1)) , las particiones no tuvieron tiempo de asignarse antes de que el programa pasara a la siguiente llamada al método (es decir seekToEnd() ).
El código de trabajo podría verse así
class Consumer(val consumer: KafkaConsumer<String, ConsumerRecord<String>>) { fun run() { // Initialization val pollDuration = 30 // seconds consumer.poll((Duration.ofSeconds(pollDuration)) // Dummy poll to get assigned partitions // Seek to end and commit new offset consumer.seekToEnd(emptyList()) consumer.commitSync() while (true) { val records = consumer.poll(Duration.ofSeconds(pollDuration)) // perform record analysis and commitSync() } } } }El método seekToEnd requiere la información sobre la partición real (en términos de Kafka TopicPartition ) en la que planea hacer que su consumidor lea desde el final.
No estoy familiarizado con la API de Kotlin, pero al revisar los JavaDocs en el método seekToEnd de KafkaConsumer , verá que solicita una colección de TopicPartitions.
Como actualmente está usando emptyList() , no tendrá ningún impacto, tal como lo observó.