diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaAdminProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaAdminProperties.java new file mode 100644 index 000000000..5cecb2d1b --- /dev/null +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaAdminProperties.java @@ -0,0 +1,62 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.kafka.properties; + +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +/** + * Properties for configuring topics. + * + * @author Gary Russell + * @since 2.0 + * + */ +public class KafkaAdminProperties { + + private Short replicationFactor; + + private Map> replicasAssignments = new HashMap<>(); + + private Map configuration = new HashMap<>(); + + public Short getReplicationFactor() { + return this.replicationFactor; + } + + public void setReplicationFactor(Short replicationFactor) { + this.replicationFactor = replicationFactor; + } + + public Map> getReplicasAssignments() { + return this.replicasAssignments; + } + + public void setReplicasAssignments(Map> replicasAssignments) { + this.replicasAssignments = replicasAssignments; + } + + public Map getConfiguration() { + return this.configuration; + } + + public void setConfiguration(Map configuration) { + this.configuration = configuration; + } + +} diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java index 5b5ea126b..5b2a6b8b8 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java @@ -87,7 +87,7 @@ public class KafkaBinderConfigurationProperties { private String requiredAcks = "1"; - private int replicationFactor = 1; + private short replicationFactor = 1; private int fetchSize = 1024 * 1024; @@ -345,11 +345,11 @@ public class KafkaBinderConfigurationProperties { this.requiredAcks = requiredAcks; } - public int getReplicationFactor() { + public short getReplicationFactor() { return this.replicationFactor; } - public void setReplicationFactor(int replicationFactor) { + public void setReplicationFactor(short replicationFactor) { this.replicationFactor = replicationFactor; } diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java index 8e00f2507..92cfed700 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java @@ -81,6 +81,8 @@ public class KafkaConsumerProperties { private Map configuration = new HashMap<>(); + private KafkaAdminProperties admin = new KafkaAdminProperties(); + public boolean isAutoCommitOffset() { return this.autoCommitOffset; } @@ -204,4 +206,12 @@ public class KafkaConsumerProperties { this.idleEventInterval = idleEventInterval; } + public KafkaAdminProperties getAdmin() { + return this.admin; + } + + public void setAdmin(KafkaAdminProperties admin) { + this.admin = admin; + } + } diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java index 104ae5834..ae35dd324 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java @@ -44,6 +44,8 @@ public class KafkaProducerProperties { private Map configuration = new HashMap<>(); + private KafkaAdminProperties admin = new KafkaAdminProperties(); + public int getBufferSize() { return this.bufferSize; } @@ -101,6 +103,15 @@ public class KafkaProducerProperties { this.configuration = configuration; } + public KafkaAdminProperties getAdmin() { + return this.admin; + } + + public void setAdmin(KafkaAdminProperties admin) { + this.admin = admin; + } + + public enum CompressionType { none, gzip, diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java index 2a7a0b0e7..a4bd6424f 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -18,6 +18,7 @@ package org.springframework.cloud.stream.binder.kafka.provisioning; import java.util.Collection; import java.util.Collections; +import java.util.List; import java.util.Map; import java.util.Set; import java.util.concurrent.Callable; @@ -44,6 +45,7 @@ import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.binder.BinderException; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaAdminProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; @@ -77,19 +79,18 @@ public class KafkaTopicProvisioner implements ProvisioningProvider adminClientProperties; private RetryOperations metadataRetryOperations; - private final int operationTimeout = DEFAULT_OPERATION_TIMEOUT; - public KafkaTopicProvisioner(KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties, KafkaProperties kafkaProperties) { Assert.isTrue(kafkaProperties != null, "KafkaProperties cannot be null"); - Map adminClientProperties = kafkaProperties.buildAdminProperties(); + this.adminClientProperties = kafkaProperties.buildAdminProperties(); this.configurationProperties = kafkaBinderConfigurationProperties; normalalizeBootPropsWithBinder(adminClientProperties, kafkaProperties, kafkaBinderConfigurationProperties); - this.adminClient = AdminClient.create(adminClientProperties); } /** @@ -118,33 +119,39 @@ public class KafkaTopicProvisioner implements ProvisioningProvider properties) { + public ProducerDestination provisionProducerDestination(final String name, + ExtendedProducerProperties properties) { + if (this.logger.isInfoEnabled()) { this.logger.info("Using kafka topic for outbound: " + name); } KafkaTopicUtils.validateTopicName(name); - createTopic(name, properties.getPartitionCount(), false); - if (this.configurationProperties.isAutoCreateTopics() && adminClient != null) { - DescribeTopicsResult describeTopicsResult = adminClient.describeTopics(Collections.singletonList(name)); - KafkaFuture> all = describeTopicsResult.all(); + try (AdminClient adminClient = AdminClient.create(this.adminClientProperties)) { + createTopic(adminClient, name, properties.getPartitionCount(), false, properties.getExtension().getAdmin()); + if (this.configurationProperties.isAutoCreateTopics()) { + DescribeTopicsResult describeTopicsResult = adminClient.describeTopics(Collections.singletonList(name)); + KafkaFuture> all = describeTopicsResult.all(); - try { - Map topicDescriptions = all.get(operationTimeout, TimeUnit.SECONDS); - TopicDescription topicDescription = topicDescriptions.get(name); - int partitions = topicDescription.partitions().size(); - return new KafkaProducerDestination(name, partitions); + try { + Map topicDescriptions = all.get(this.operationTimeout, TimeUnit.SECONDS); + TopicDescription topicDescription = topicDescriptions.get(name); + int partitions = topicDescription.partitions().size(); + return new KafkaProducerDestination(name, partitions); + } + catch (Exception e) { + throw new ProvisioningException("Problems encountered with partitions finding", e); + } } - catch (Exception e) { - throw new ProvisioningException("Problems encountered with partitions finding", e); + else { + return new KafkaProducerDestination(name); } } - else { - return new KafkaProducerDestination(name); - } } @Override - public ConsumerDestination provisionConsumerDestination(final String name, final String group, ExtendedConsumerProperties properties) { + public ConsumerDestination provisionConsumerDestination(final String name, final String group, + ExtendedConsumerProperties properties) { + KafkaTopicUtils.validateTopicName(name); boolean anonymous = !StringUtils.hasText(group); Assert.isTrue(!anonymous || !properties.getExtension().isEnableDlq(), @@ -153,27 +160,35 @@ public class KafkaTopicProvisioner implements ProvisioningProvider> all = describeTopicsResult.all(); - try { - Map topicDescriptions = all.get(operationTimeout, TimeUnit.SECONDS); - TopicDescription topicDescription = topicDescriptions.get(name); - int partitions = topicDescription.partitions().size(); - ConsumerDestination dlqTopic = createDlqIfNeedBe(name, group, properties, anonymous, partitions); - if (dlqTopic != null) { - return dlqTopic; + try (AdminClient adminClient = createAdminClient()) { + createTopic(adminClient, name, partitionCount, properties.getExtension().isAutoRebalanceEnabled(), + properties.getExtension().getAdmin()); + if (this.configurationProperties.isAutoCreateTopics()) { + DescribeTopicsResult describeTopicsResult = adminClient.describeTopics(Collections.singletonList(name)); + KafkaFuture> all = describeTopicsResult.all(); + try { + Map topicDescriptions = all.get(operationTimeout, TimeUnit.SECONDS); + TopicDescription topicDescription = topicDescriptions.get(name); + int partitions = topicDescription.partitions().size(); + ConsumerDestination dlqTopic = createDlqIfNeedBe(adminClient, name, group, properties, anonymous, + partitions); + if (dlqTopic != null) { + return dlqTopic; + } + return new KafkaConsumerDestination(name, partitions); + } + catch (Exception e) { + throw new ProvisioningException("provisioning exception", e); } - return new KafkaConsumerDestination(name, partitions); - } - catch (Exception e) { - throw new ProvisioningException("provisioning exception", e); } } return new KafkaConsumerDestination(name); } + AdminClient createAdminClient() { + return AdminClient.create(this.adminClientProperties); + } + /** * In general, binder properties supersede boot kafka properties. * The one exception is the bootstrap servers. In that case, we should only override @@ -209,14 +224,15 @@ public class KafkaTopicProvisioner implements ProvisioningProvider properties, boolean anonymous, int partitions) { if (properties.getExtension().isEnableDlq() && !anonymous) { String dlqTopic = StringUtils.hasText(properties.getExtension().getDlqName()) ? properties.getExtension().getDlqName() : "error." + name + "." + group; try { - createTopicAndPartitions(dlqTopic, partitions, properties.getExtension().isAutoRebalanceEnabled()); + createTopicAndPartitions(adminClient, dlqTopic, partitions, + properties.getExtension().isAutoRebalanceEnabled(), properties.getExtension().getAdmin()); } catch (Throwable throwable) { if (throwable instanceof Error) { @@ -231,9 +247,10 @@ public class KafkaTopicProvisioner implements ProvisioningProvider> namesFutures = listTopicsResult.names(); @@ -298,14 +316,26 @@ public class KafkaTopicProvisioner implements ProvisioningProvider { - NewTopic newTopic = new NewTopic(topicName, effectivePartitionCount, - (short) configurationProperties.getReplicationFactor()); + NewTopic newTopic; + Map> replicasAssignments = adminProperties.getReplicasAssignments(); + if (replicasAssignments != null && replicasAssignments.size() > 0) { + newTopic = new NewTopic(topicName, adminProperties.getReplicasAssignments()); + } + else { + newTopic = new NewTopic(topicName, effectivePartitionCount, + adminProperties.getReplicationFactor() != null + ? adminProperties.getReplicationFactor() + : configurationProperties.getReplicationFactor()); + } + if (adminProperties.getConfiguration().size() > 0) { + newTopic.configs(adminProperties.getConfiguration()); + } CreateTopicsResult createTopicsResult = adminClient.createTopics(Collections.singletonList(newTopic)); try { createTopicsResult.all().get(operationTimeout, TimeUnit.SECONDS); @@ -318,6 +348,10 @@ public class KafkaTopicProvisioner implements ProvisioningProvider