From c389fa3fa442632ff512946d04db6e46588a1f74 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 24 Jan 2017 16:29:10 -0500 Subject: [PATCH] Ignore 'TopicExistsException' on creation Fixes a race condtion when topics are created in a stream and a TopicExistsException is thrown. Fixes #83 Explicitly checking for TopicExistsException Cleanup --- .../kafka/KafkaMessageChannelBinder.java | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java index 241ee8501..6a05f67ff 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaMessageChannelBinder.java @@ -470,8 +470,22 @@ public class KafkaMessageChannelBinder extends @Override public Object doWithRetry(RetryContext context) throws RuntimeException { - adminUtilsOperation.invokeCreateTopic(zkUtils, topicName, effectivePartitionCount, - configurationProperties.getReplicationFactor(), new Properties()); + try { + adminUtilsOperation.invokeCreateTopic(zkUtils, topicName, effectivePartitionCount, + configurationProperties.getReplicationFactor(), new Properties()); + } + catch (Exception e) { + String exceptionClass = e.getClass().getName(); + if (exceptionClass.equals("kafka.common.TopicExistsException") || + exceptionClass.equals("org.apache.kafka.common.errors.TopicExistsException")){ + if (logger.isWarnEnabled()) { + logger.warn("Attempt to create topic: " + topicName + ". Topic already exists."); + } + } + else { + throw e; + } + } return null; } });