He estado tratando de hacer un trabajo de POC para Spring Kafka. Específicamente, quería experimentar cuáles son las mejores prácticas en términos de manejo de errores al consumir mensajes dentro de Kafka.
Me pregunto si alguien puede ayudar con:
El ejemplo de código para 2 se da a continuación:
Dado que AckMode está configurado en RECORD, que según la documentación :
comete el desplazamiento cuando el oyente regresa después de procesar el registro.
Hubiera pensado que el desplazamiento no se incrementaría si el método de escucha arrojara una excepción. Sin embargo, este no fue el caso cuando lo probé usando la combinación de código/configuración/comando a continuación. El desplazamiento aún se actualiza y el siguiente mensaje continúa procesándose.
Mi configuración:
private Map<String, Object> producerConfigs() { Map<String, Object> props = new HashMap<>(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.0.1:9092"); props.put(ProducerConfig.RETRIES_CONFIG, 0); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.LINGER_MS_CONFIG, 1); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class); return props; } @Bean ConcurrentKafkaListenerContainerFactory<Integer, String> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<Integer, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(new DefaultKafkaConsumerFactory<>(consumerConfigs())); factory.getContainerProperties().setAckMode(AbstractMessageListenerContainer.AckMode.RECORD); return factory; }Mi código:
@Component public class KafkaMessageListener{ @KafkaListener(topicPartitions = {@TopicPartition( topic = "my-replicated-topic", partitionOffsets = @PartitionOffset(partition = "0", initialOffset = "0", relativeToCurrent = "true"))}) public void onReplicatedTopicMessage(ConsumerRecord<Integer, String> data) throws InterruptedException { throw new RuntimeException("Oops!"); }Comando para verificar el desplazamiento:
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group test-groupEstoy usando kafka_2.12-0.10.2.0 y org.springframework.kafka:spring-kafka:1.1.3.RELEASE
El contenedor (a través de ContainerProperties ) tiene una propiedad, ackOnError , que es verdadera de forma predeterminada...
/** * Set whether or not the container should commit offsets (ack messages) where the * listener throws exceptions. This works in conjunction with {@link #ackMode} and is * effective only when the kafka property {@code enable.auto.commit} is {@code false}; * it is not applicable to manual ack modes. When this property is set to {@code true} * (the default), all messages handled will have their offset committed. When set to * {@code false}, offsets will be committed only for successfully handled messages. * Manual acks will be always be applied. Bear in mind that, if the next message is * successfully handled, its offset will be committed, effectively committing the * offset of the failed message anyway, so this option has limited applicability. * Perhaps useful for a component that starts throwing exceptions consistently; * allowing it to resume when restarted from the last successfully processed message. * @param ackOnError whether the container should acknowledge messages that throw * exceptions. */ public void setAckOnError(boolean ackOnError) { this.ackOnError = ackOnError; }Tenga en cuenta, sin embargo, que si el siguiente mensaje es exitoso, su compensación se confirmará de todos modos, lo que efectivamente también compromete la compensación fallida.
EDITAR
A partir de la versión 2.3, ackOnError ahora es false de forma predeterminada.