He escrito un programa que crea temas de Kafka al iniciarse con AdminClient. Si existe un tema, se lanza una TopicExistsException. Sin embargo, no soy capaz de atraparlo. Código relacionado:
try { CreateTopicsResult result = adminClient.createTopics(Collections.singleton(new NewTopic(topic, 2, (short) 1))); result.all().get(5, TimeUnit.SECONDS); } catch (TopicExistsException e){ log.info("topic exists: {}", topic); } catch ( InterruptedException | TimeoutException | ExecutionException e){ log.error("Topic create failed with: {}", e.getMessage()); }El comportamiento esperado sería que se ejecute la primera captura, pero según los registros, esto no es lo que sucede:
2020-12-01 13:18:56,462 ERROR [restartedMain] com.mypackage.KafkaService@lambda$configureTopics$0:75 - Topic create failed with: org.apache.kafka.common.errors.TopicExistsException: Topic 'MEASUREMENT_TEST' already exists.Mis importaciones son para las excepciones:
import org.apache.kafka.common.errors.TopicExistsException; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeoutException; import java.lang.InterruptedException;¿Cualquier sugerencia? ¿Hay algo que me estoy perdiendo?
Esto se debe a que la excepción de que existe el tema se envuelve en la ExecutionException que se genera cuando se completa KafkaFuture .
Puede ver esto registrando el seguimiento completo de la pila, no solo la parte getMessage() :
java.util.concurrent.ExecutionException: org.apache.kafka.common.errors.TopicExistsException: Topic 'hello' already exists. at org.apache.kafka.common.internals.KafkaFutureImpl.wrapAndThrow(KafkaFutureImpl.java:45) ~[kafka-clients-2.0.0.jar:?] at org.apache.kafka.common.internals.KafkaFutureImpl.access$000(KafkaFutureImpl.java:32) ~[kafka-clients-2.0.0.jar:?] at org.apache.kafka.common.internals.KafkaFutureImpl$SingleWaiter.await(KafkaFutureImpl.java:89) ~[kafka-clients-2.0.0.jar:?] at org.apache.kafka.common.internals.KafkaFutureImpl.get(KafkaFutureImpl.java:262) ~[kafka-clients-2.0.0.jar:?] Caused by: org.apache.kafka.common.errors.TopicExistsException: Topic 'hello' already exists. Entonces, si desea capturar TopicExistsException , debe hacer algo como:
try { CreateTopicsResult result = adminClient.createTopics(Collections.singleton(new NewTopic(topic, 2, (short) 1))); result.all().get(5, TimeUnit.SECONDS); } catch (ExecutionException e){ if(e.getCause() != null && e.getCause() instanceof TopicExistsException) { log.info("topic exists: {}", topic); } else { log.error("Topic create failed with: {}", e); } } catch ( InterruptedException | TimeoutException){ log.error("Topic create failed with: {}", e.getMessage()); }