diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index 48b9e339f..d3d5d4b32 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -203,8 +203,8 @@ public class KafkaTopicProvisioner implements } } private ConsumerDestination doProvisionConsumerDestination(final String name, - final String group, - ExtendedConsumerProperties properties) { + final String group, + ExtendedConsumerProperties properties) { final KafkaConsumerDestination kafkaConsumerDestination = new KafkaConsumerDestination(name); if (properties.getExtension().isDestinationIsPattern()) { Assert.isTrue(!properties.getExtension().isEnableDlq(), @@ -262,7 +262,8 @@ public class KafkaTopicProvisioner implements .describeTopics(Collections.singletonList(topicName)); KafkaFuture> all = describeTopicsResult.all(); return all.get(this.operationTimeout, TimeUnit.SECONDS); - } catch (Exception ex) { + } + catch (Exception ex) { throw new ProvisioningException("Problems encountered with partitions finding for: " + topicName, ex); } });