Business
Jobs
  • About Us
  • Solutions
    • Job Postings
      Post your job and receive qualified candidates in 48h.
    • Candidate Assessments
      500+ technical and psychological tests, plus anti-fraud.
    • Headhunting
      Tailor-made executive search from start to finish.
    • Payroll + EOR
      Payroll dispersal and EOR across 15+ LATAM countries.
  • Pricing
  • Jobs

0

350
Views
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 answers
Answer question

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 Report

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 Report
Answer question
Find remote jobs

Discover the new way to find a job!

Top jobs
Top job categories
Business
Post vacancy Pricing Sales
Legal
Terms and conditions Privacy policy
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Show me some job opportunities
There's an error!