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 d3d5d4b32..402c248bf 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 @@ -561,8 +561,7 @@ public class KafkaTopicProvisioner implements // In some cases, the above partition query may not throw an UnknownTopic..Exception for various reasons. // For that, we are forcing another query to ensure that the topic is present on the server. if (CollectionUtils.isEmpty(partitions)) { - try (AdminClient adminClient = AdminClient - .create(this.adminClientProperties)) { + try (AdminClient adminClient = createAdminClient()) { final DescribeTopicsResult describeTopicsResult = adminClient .describeTopics(Collections.singletonList(topicName));