Empresas
Empleos
  • Sobre nosotros
  • Soluciones
    • Publicación de vacantes
      Publica tu vacante y recibe candidatos calificados en 48h.
    • Evaluación de candidatos
      500+ pruebas técnicas y psicológicas, más anti-fraude.
    • Headhunting
      Búsqueda ejecutiva a la medida de principio a fin.
    • Nómina + EOR
      Dispersión de nómina y EOR en más de 15 países de LATAM.
  • Precios
  • Empleos

0

355
Vistas
KafkaConsumer: `seekToEnd()` no hace que el consumidor consuma desde la última compensación

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 }
over 4 years ago · Santiago Trujillo
2 Respuestas
Responde la pregunta

0

¡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() } } } }
over 4 years ago · Santiago Trujillo Denunciar

0

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ó.

over 4 years ago · Santiago Trujillo Denunciar
Responde la pregunta
Encuentra empleos remotos

¡Descubre la nueva forma de encontrar empleo!

Top de empleos
Top categorías de empleo
Empresas
Publicar vacante Precios Comercial
Legal
Términos y condiciones Política de privacidad
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomiéndame algunas ofertas
Necesito ayuda