diff --git a/pom.xml b/pom.xml index 72cacd09c..754452b72 100644 --- a/pom.xml +++ b/pom.xml @@ -23,6 +23,7 @@ spring-cloud-stream-binder-kafka-docs spring-cloud-stream-binder-kafka-0.9-test spring-cloud-stream-binder-kafka-0.10.0-test + spring-cloud-stream-binder-kafka-core diff --git a/spring-cloud-stream-binder-kafka-0.10.0-test/pom.xml b/spring-cloud-stream-binder-kafka-0.10.0-test/pom.xml index 61a235e25..56a60c4cc 100644 --- a/spring-cloud-stream-binder-kafka-0.10.0-test/pom.xml +++ b/spring-cloud-stream-binder-kafka-0.10.0-test/pom.xml @@ -25,6 +25,11 @@ + + org.springframework.cloud + spring-cloud-stream-binder-kafka-core + 1.2.0.BUILD-SNAPSHOT + org.springframework.cloud spring-cloud-stream-binder-kafka diff --git a/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java b/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java index 31ca93e5b..adb53e52a 100644 --- a/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java +++ b/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java @@ -16,8 +16,10 @@ package org.springframework.cloud.stream.binder.kafka; +import org.springframework.cloud.stream.binder.kafka.admin.AdminUtilsOperation; import org.springframework.cloud.stream.binder.kafka.admin.Kafka09AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.context.support.GenericApplicationContext; import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; @@ -35,14 +37,18 @@ public class Kafka09TestBinder extends AbstractKafkaTestBinder { public Kafka09TestBinder(KafkaBinderConfigurationProperties binderConfiguration) { try { - KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(binderConfiguration); + AdminUtilsOperation adminUtilsOperation = new Kafka09AdminUtilsOperation(); + KafkaTopicProvisioner provisioningProvider = + new KafkaTopicProvisioner(binderConfiguration, adminUtilsOperation); + provisioningProvider.afterPropertiesSet(); + + KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(binderConfiguration, provisioningProvider); binder.setCodec(getCodec()); ProducerListener producerListener = new LoggingProducerListener(); binder.setProducerListener(producerListener); GenericApplicationContext context = new GenericApplicationContext(); context.refresh(); binder.setApplicationContext(context); - binder.setAdminUtilsOperation(new Kafka09AdminUtilsOperation()); binder.afterPropertiesSet(); this.setBinder(binder); } diff --git a/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka_09_BinderTests.java b/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka_09_BinderTests.java index 583387877..e9c470ce7 100644 --- a/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka_09_BinderTests.java +++ b/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka_09_BinderTests.java @@ -34,13 +34,12 @@ import org.junit.ClassRule; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.Spy; import org.springframework.cloud.stream.binder.kafka.admin.Kafka09AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.test.core.BrokerAddress; import org.springframework.kafka.test.rule.KafkaEmbedded; -import org.springframework.retry.RetryOperations; /** * Integration tests for the {@link KafkaMessageChannelBinder}. @@ -93,11 +92,6 @@ public class Kafka_09_BinderTests extends KafkaBinderTests { return consumerFactory().createConsumer().partitionsFor(topic).size(); } - @Override - protected void setMetadataRetryOperations(Binder binder, RetryOperations retryOperations) { - ((Kafka09TestBinder) binder).getBinder().setMetadataRetryOperations(retryOperations); - } - @Override protected ZkUtils getZkUtils(KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties) { final ZkClient zkClient = new ZkClient(kafkaBinderConfigurationProperties.getZkConnectionString(), diff --git a/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafka09BinderTests.java b/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafka09BinderTests.java index 2f058733a..eb386248d 100644 --- a/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafka09BinderTests.java +++ b/spring-cloud-stream-binder-kafka-0.9-test/src/test/java/org/springframework/cloud/stream/binder/kafka/RawModeKafka09BinderTests.java @@ -24,6 +24,8 @@ import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.HeaderMode; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; diff --git a/spring-cloud-stream-binder-kafka-core/pom.xml b/spring-cloud-stream-binder-kafka-core/pom.xml new file mode 100644 index 000000000..dd863135a --- /dev/null +++ b/spring-cloud-stream-binder-kafka-core/pom.xml @@ -0,0 +1,46 @@ + + + 4.0.0 + + org.springframework.cloud + spring-cloud-stream-binder-kafka-parent + 1.2.0.BUILD-SNAPSHOT + + spring-cloud-stream-binder-kafka-core + Spring Cloud Stream Kafka Binder Core + http://projects.spring.io/spring-cloud + + Pivotal Software, Inc. + http://www.spring.io + + + + + + + + org.springframework.cloud + spring-cloud-stream + + + org.apache.kafka + kafka-clients + + + org.apache.kafka + kafka_2.11 + + + org.springframework.integration + spring-integration-kafka + ${spring-integration-kafka.version} + + + org.apache.avro + avro-compiler + + + + + + diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/AdminUtilsOperation.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/AdminUtilsOperation.java similarity index 100% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/AdminUtilsOperation.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/AdminUtilsOperation.java diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka09AdminUtilsOperation.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka09AdminUtilsOperation.java similarity index 99% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka09AdminUtilsOperation.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka09AdminUtilsOperation.java index dfacbf3bc..0723f1ed6 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka09AdminUtilsOperation.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka09AdminUtilsOperation.java @@ -142,4 +142,4 @@ public class Kafka09AdminUtilsOperation implements AdminUtilsOperation { ReflectionUtils.handleReflectionException(e); } } -} +} \ No newline at end of file diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka10AdminUtilsOperation.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka10AdminUtilsOperation.java similarity index 100% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka10AdminUtilsOperation.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/admin/Kafka10AdminUtilsOperation.java diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/JaasLoginModuleConfiguration.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/JaasLoginModuleConfiguration.java similarity index 97% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/JaasLoginModuleConfiguration.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/JaasLoginModuleConfiguration.java index fc2e97c21..33e9865f1 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/JaasLoginModuleConfiguration.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/JaasLoginModuleConfiguration.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka.config; +package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java similarity index 98% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java index 158073a6d..4bb30d952 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBinderConfigurationProperties.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka.config; +package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBindingProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBindingProperties.java similarity index 94% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBindingProperties.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBindingProperties.java index 6bd233d63..e7a5d731d 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBindingProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaBindingProperties.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka; +package org.springframework.cloud.stream.binder.kafka.properties; /** * @author Marius Bogoevici diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java similarity index 97% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java index e824e356f..fec952db7 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka; +package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaExtendedBindingProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaExtendedBindingProperties.java similarity index 96% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaExtendedBindingProperties.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaExtendedBindingProperties.java index 25be14982..5b1cc9117 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaExtendedBindingProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaExtendedBindingProperties.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka; +package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaProducerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java similarity index 96% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaProducerProperties.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java index 2568a1c51..df0c78ebb 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaProducerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaProducerProperties.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka; +package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; 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 new file mode 100644 index 000000000..99c0b2fed --- /dev/null +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/provisioning/KafkaTopicProvisioner.java @@ -0,0 +1,310 @@ +/* + * Copyright 2014-2016 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.provisioning; + +import java.util.Collection; +import java.util.Properties; +import java.util.concurrent.Callable; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.apache.kafka.common.PartitionInfo; +import org.apache.kafka.common.security.JaasUtils; + +import org.springframework.beans.factory.InitializingBean; +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.admin.AdminUtilsOperation; +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; +import org.springframework.cloud.stream.binder.kafka.utils.KafkaTopicUtils; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.cloud.stream.provisioning.ProducerDestination; +import org.springframework.cloud.stream.provisioning.ProvisioningException; +import org.springframework.cloud.stream.provisioning.ProvisioningProvider; +import org.springframework.retry.RetryCallback; +import org.springframework.retry.RetryContext; +import org.springframework.retry.RetryOperations; +import org.springframework.retry.backoff.ExponentialBackOffPolicy; +import org.springframework.retry.policy.SimpleRetryPolicy; +import org.springframework.retry.support.RetryTemplate; +import org.springframework.util.Assert; +import org.springframework.util.StringUtils; + +import kafka.common.ErrorMapping; +import kafka.utils.ZkUtils; + +/** + * Kafka implementation for {@link ProvisioningProvider} + * + * @author Soby Chacko + */ +public class KafkaTopicProvisioner implements ProvisioningProvider, + ExtendedProducerProperties>, InitializingBean { + + private final Log logger = LogFactory.getLog(getClass()); + + private final KafkaBinderConfigurationProperties configurationProperties; + + private final AdminUtilsOperation adminUtilsOperation; + + private RetryOperations metadataRetryOperations; + + public KafkaTopicProvisioner(KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties, + AdminUtilsOperation adminUtilsOperation) { + this.configurationProperties = kafkaBinderConfigurationProperties; + this.adminUtilsOperation = adminUtilsOperation; + } + + /** + * + * @param metadataRetryOperations the retry configuration + */ + public void setMetadataRetryOperations(RetryOperations metadataRetryOperations) { + this.metadataRetryOperations = metadataRetryOperations; + } + + @Override + public void afterPropertiesSet() throws Exception { + if (this.metadataRetryOperations == null) { + RetryTemplate retryTemplate = new RetryTemplate(); + + SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy(); + simpleRetryPolicy.setMaxAttempts(10); + retryTemplate.setRetryPolicy(simpleRetryPolicy); + + ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy(); + backOffPolicy.setInitialInterval(100); + backOffPolicy.setMultiplier(2); + backOffPolicy.setMaxInterval(1000); + retryTemplate.setBackOffPolicy(backOffPolicy); + this.metadataRetryOperations = retryTemplate; + } + } + + @Override + public ProducerDestination provisionProducerDestination(final String name, ExtendedProducerProperties properties) { + if (this.logger.isInfoEnabled()) { + this.logger.info("Using kafka topic for outbound: " + name); + } + KafkaTopicUtils.validateTopicName(name); + createTopicsIfAutoCreateEnabledAndAdminUtilsPresent(name, properties.getPartitionCount()); + if (this.configurationProperties.isAutoCreateTopics() && adminUtilsOperation != null) { + final ZkUtils zkUtils = ZkUtils.apply(this.configurationProperties.getZkConnectionString(), + this.configurationProperties.getZkSessionTimeout(), + this.configurationProperties.getZkConnectionTimeout(), + JaasUtils.isZkSecurityEnabled()); + int partitions = adminUtilsOperation.partitionSize(name, zkUtils); + return new KafkaProducerDestination(name, partitions); + } + else { + return new KafkaProducerDestination(name); + } + } + + @Override + 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(), + "DLQ support is not available for anonymous subscriptions"); + if (properties.getInstanceCount() == 0) { + throw new IllegalArgumentException("Instance count cannot be zero"); + } + int partitionCount = properties.getInstanceCount() * properties.getConcurrency(); + createTopicsIfAutoCreateEnabledAndAdminUtilsPresent(name, partitionCount); + if (this.configurationProperties.isAutoCreateTopics() && adminUtilsOperation != null) { + final ZkUtils zkUtils = ZkUtils.apply(this.configurationProperties.getZkConnectionString(), + this.configurationProperties.getZkSessionTimeout(), + this.configurationProperties.getZkConnectionTimeout(), + JaasUtils.isZkSecurityEnabled()); + int partitions = adminUtilsOperation.partitionSize(name, zkUtils); + if (properties.getExtension().isEnableDlq() && !anonymous) { + String dlqTopic = "error." + name + "." + group; + createTopicAndPartitions(dlqTopic, partitions); + return new KafkaConsumerDestination(name, partitions, dlqTopic); + } + return new KafkaConsumerDestination(name, partitions); + } + return new KafkaConsumerDestination(name); + } + + private void createTopicsIfAutoCreateEnabledAndAdminUtilsPresent(final String topicName, final int partitionCount) { + if (this.configurationProperties.isAutoCreateTopics() && adminUtilsOperation != null) { + createTopicAndPartitions(topicName, partitionCount); + } + else if (this.configurationProperties.isAutoCreateTopics() && adminUtilsOperation == null) { + this.logger.warn("Auto creation of topics is enabled, but Kafka AdminUtils class is not present on the classpath. " + + "No topic will be created by the binder"); + } + else if (!this.configurationProperties.isAutoCreateTopics()) { + this.logger.info("Auto creation of topics is disabled."); + } + } + + /** + * Creates a Kafka topic if needed, or try to increase its partition count to the + * desired number. + */ + private void createTopicAndPartitions(final String topicName, final int partitionCount) { + final ZkUtils zkUtils = ZkUtils.apply(this.configurationProperties.getZkConnectionString(), + this.configurationProperties.getZkSessionTimeout(), + this.configurationProperties.getZkConnectionTimeout(), + JaasUtils.isZkSecurityEnabled()); + try { + short errorCode = adminUtilsOperation.errorCodeFromTopicMetadata(topicName, zkUtils); + if (errorCode == ErrorMapping.NoError()) { + // only consider minPartitionCount for resizing if autoAddPartitions is true + int effectivePartitionCount = this.configurationProperties.isAutoAddPartitions() + ? Math.max(this.configurationProperties.getMinPartitionCount(), partitionCount) + : partitionCount; + int partitionSize = adminUtilsOperation.partitionSize(topicName, zkUtils); + + if (partitionSize < effectivePartitionCount) { + if (this.configurationProperties.isAutoAddPartitions()) { + adminUtilsOperation.invokeAddPartitions(zkUtils, topicName, effectivePartitionCount, null, false); + } + else { + throw new ProvisioningException("The number of expected partitions was: " + partitionCount + ", but " + + partitionSize + (partitionSize > 1 ? " have " : " has ") + "been found instead." + + "Consider either increasing the partition count of the topic or enabling " + + "`autoAddPartitions`"); + } + } + } + else if (errorCode == ErrorMapping.UnknownTopicOrPartitionCode()) { + // always consider minPartitionCount for topic creation + final int effectivePartitionCount = Math.max(this.configurationProperties.getMinPartitionCount(), + partitionCount); + + this.metadataRetryOperations.execute(new RetryCallback() { + + @Override + public Object doWithRetry(RetryContext context) throws RuntimeException { + + adminUtilsOperation.invokeCreateTopic(zkUtils, topicName, effectivePartitionCount, + configurationProperties.getReplicationFactor(), new Properties()); + return null; + } + }); + } + else { + throw new ProvisioningException("Error fetching Kafka topic metadata: ", + ErrorMapping.exceptionFor(errorCode)); + } + } + finally { + zkUtils.close(); + } + } + + public Collection getPartitionsForTopic(final int partitionCount, final Callable> callable) { + try { + return this.metadataRetryOperations + .execute(new RetryCallback, Exception>() { + + @Override + public Collection doWithRetry(RetryContext context) throws Exception { + Collection partitions = callable.call(); + // do a sanity check on the partition set + if (partitions.size() < partitionCount) { + throw new IllegalStateException("The number of expected partitions was: " + + partitionCount + ", but " + partitions.size() + + (partitions.size() > 1 ? " have " : " has ") + "been found instead"); + } + return partitions; + } + }); + } + catch (Exception e) { + this.logger.error("Cannot initialize Binder", e); + throw new BinderException("Cannot initialize binder:", e); + } + } + + private final class KafkaProducerDestination implements ProducerDestination { + + private final String producerDestinationName; + private final int partitions; + + private KafkaProducerDestination(String destinationName, Integer partitions) { + this.producerDestinationName = destinationName; + this.partitions = partitions; + } + + private KafkaProducerDestination(String destinationName) { + this(destinationName, 0); + } + + @Override + public String getName() { + return producerDestinationName; + } + + @Override + public String getNameForPartition(int partition) { + return producerDestinationName; + } + + @Override + public String toString() { + return "KafkaProducerDestination{" + + "producerDestinationName='" + producerDestinationName + '\'' + + ", partitions=" + partitions + + '}'; + } + } + + private final class KafkaConsumerDestination implements ConsumerDestination { + + private final String consumerDestinationName; + private final int partitions; + private final String dlqName; + + private KafkaConsumerDestination(String consumerDestinationName) { + this(consumerDestinationName, 0, null); + } + + private KafkaConsumerDestination(String consumerDestinationName, int partitions) { + this(consumerDestinationName, partitions, null); + } + + private KafkaConsumerDestination(String consumerDestinationName, Integer partitions, String dlqName) { + this.consumerDestinationName = consumerDestinationName; + this.partitions = partitions; + this.dlqName = dlqName; + } + + @Override + public String getName() { + return this.consumerDestinationName; + } + + @Override + public String toString() { + return "KafkaConsumerDestination{" + + "consumerDestinationName='" + consumerDestinationName + '\'' + + ", partitions=" + partitions + + ", dlqName='" + dlqName + '\'' + + '}'; + } + + } + +} diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaTopicUtils.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/utils/KafkaTopicUtils.java similarity index 95% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaTopicUtils.java rename to spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/utils/KafkaTopicUtils.java index b7434ddbe..8b937bfbe 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaTopicUtils.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/utils/KafkaTopicUtils.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka; +package org.springframework.cloud.stream.binder.kafka.utils; import java.io.UnsupportedEncodingException; diff --git a/spring-cloud-stream-binder-kafka/pom.xml b/spring-cloud-stream-binder-kafka/pom.xml index deafb7966..71f427a38 100644 --- a/spring-cloud-stream-binder-kafka/pom.xml +++ b/spring-cloud-stream-binder-kafka/pom.xml @@ -14,6 +14,11 @@ + + org.springframework.cloud + spring-cloud-stream-binder-kafka-core + 1.2.0.BUILD-SNAPSHOT + org.springframework.boot spring-boot-configuration-processor diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java index d36552ec1..e1e89946c 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderHealthIndicator.java @@ -29,7 +29,7 @@ import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.springframework.boot.actuate.health.Health; import org.springframework.boot.actuate.health.HealthIndicator; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; /** * Health indicator for Kafka. diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderJaasInitializerListener.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderJaasInitializerListener.java index c5e649825..4d06df006 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderJaasInitializerListener.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderJaasInitializerListener.java @@ -28,7 +28,7 @@ import org.apache.kafka.common.security.JaasUtils; import org.springframework.beans.BeansException; import org.springframework.beans.factory.DisposableBean; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; import org.springframework.context.ApplicationListener; 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 6a05f67ff..ae54100be 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 @@ -22,34 +22,30 @@ import java.util.Arrays; import java.util.Collection; import java.util.HashMap; import java.util.Map; -import java.util.Properties; import java.util.UUID; +import java.util.concurrent.Callable; -import kafka.common.ErrorMapping; -import kafka.utils.ZkUtils; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRecord; -import org.apache.kafka.clients.producer.Callback; -import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.clients.producer.ProducerConfig; -import org.apache.kafka.clients.producer.ProducerRecord; -import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.PartitionInfo; -import org.apache.kafka.common.security.JaasUtils; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.apache.kafka.common.utils.Utils; -import org.springframework.beans.factory.DisposableBean; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; import org.springframework.cloud.stream.binder.Binder; -import org.springframework.cloud.stream.binder.BinderException; import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; -import org.springframework.cloud.stream.binder.kafka.admin.AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +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.KafkaExtendedBindingProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; +import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.cloud.stream.provisioning.ProducerDestination; import org.springframework.context.Lifecycle; import org.springframework.expression.common.LiteralExpression; import org.springframework.expression.spel.standard.SpelExpressionParser; @@ -65,19 +61,17 @@ import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.kafka.listener.ErrorHandler; import org.springframework.kafka.listener.config.ContainerProperties; import org.springframework.kafka.support.ProducerListener; +import org.springframework.kafka.support.SendResult; import org.springframework.kafka.support.TopicPartitionInitialOffset; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHandler; -import org.springframework.retry.RetryCallback; -import org.springframework.retry.RetryContext; -import org.springframework.retry.RetryOperations; -import org.springframework.retry.backoff.ExponentialBackOffPolicy; -import org.springframework.retry.policy.SimpleRetryPolicy; import org.springframework.retry.support.RetryTemplate; import org.springframework.util.Assert; import org.springframework.util.CollectionUtils; import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; +import org.springframework.util.concurrent.ListenableFuture; +import org.springframework.util.concurrent.ListenableFutureCallback; /** * A {@link Binder} that uses Kafka as the underlying middleware. @@ -92,26 +86,20 @@ import org.springframework.util.StringUtils; */ public class KafkaMessageChannelBinder extends AbstractMessageChannelBinder, - ExtendedProducerProperties, Collection, String> - implements ExtendedPropertiesBinder, - DisposableBean { + ExtendedProducerProperties, KafkaTopicProvisioner> + implements ExtendedPropertiesBinder { private final KafkaBinderConfigurationProperties configurationProperties; - private RetryOperations metadataRetryOperations; - - private final Map> topicsInUse = new HashMap<>(); - private ProducerListener producerListener; - private volatile Producer dlqProducer; - private KafkaExtendedBindingProperties extendedBindingProperties = new KafkaExtendedBindingProperties(); - private AdminUtilsOperation adminUtilsOperation; + private final Map> topicsInUse = new HashMap<>(); - public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties) { - super(false, headersToMap(configurationProperties)); + public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, + KafkaTopicProvisioner provisioningProvider) { + super(false, headersToMap(configurationProperties), provisioningProvider); this.configurationProperties = configurationProperties; } @@ -131,50 +119,10 @@ public class KafkaMessageChannelBinder extends return headersToMap; } - public void setAdminUtilsOperation(AdminUtilsOperation adminUtilsOperation) { - this.adminUtilsOperation = adminUtilsOperation; - } - - /** - * Retry configuration for operations such as validating topic creation - * - * @param metadataRetryOperations the retry configuration - */ - public void setMetadataRetryOperations(RetryOperations metadataRetryOperations) { - this.metadataRetryOperations = metadataRetryOperations; - } - public void setExtendedBindingProperties(KafkaExtendedBindingProperties extendedBindingProperties) { this.extendedBindingProperties = extendedBindingProperties; } - @Override - public void onInit() throws Exception { - - if (this.metadataRetryOperations == null) { - RetryTemplate retryTemplate = new RetryTemplate(); - - SimpleRetryPolicy simpleRetryPolicy = new SimpleRetryPolicy(); - simpleRetryPolicy.setMaxAttempts(10); - retryTemplate.setRetryPolicy(simpleRetryPolicy); - - ExponentialBackOffPolicy backOffPolicy = new ExponentialBackOffPolicy(); - backOffPolicy.setInitialInterval(100); - backOffPolicy.setMultiplier(2); - backOffPolicy.setMaxInterval(1000); - retryTemplate.setBackOffPolicy(backOffPolicy); - this.metadataRetryOperations = retryTemplate; - } - } - - @Override - public void destroy() throws Exception { - if (this.dlqProducer != null) { - this.dlqProducer.close(); - this.dlqProducer = null; - } - } - public void setProducerListener(ProducerListener producerListener) { this.producerListener = producerListener; } @@ -194,49 +142,30 @@ public class KafkaMessageChannelBinder extends } @Override - protected MessageHandler createProducerMessageHandler(final String destination, + protected MessageHandler createProducerMessageHandler(final ProducerDestination destination, ExtendedProducerProperties producerProperties) throws Exception { - - KafkaTopicUtils.validateTopicName(destination); - createTopicsIfAutoCreateEnabledAndAdminUtilsPresent(destination, producerProperties.getPartitionCount()); - Collection partitions = getPartitionsForTopic(destination, producerProperties.getPartitionCount()); - + final DefaultKafkaProducerFactory producerFB = getProducerFactory(producerProperties); + Collection partitions = provisioningProvider.getPartitionsForTopic(producerProperties.getPartitionCount(), + new Callable>() { + @Override + public Collection call() throws Exception { + return producerFB.createProducer().partitionsFor(destination.getName()); + } + }); + this.topicsInUse.put(destination.getName(), partitions); if (producerProperties.getPartitionCount() < partitions.size()) { if (this.logger.isInfoEnabled()) { - this.logger.info("The `partitionCount` of the producer for topic " + destination + " is " + this.logger.info("The `partitionCount` of the producer for topic " + destination.getName() + " is " + producerProperties.getPartitionCount() + ", smaller than the actual partition count of " + partitions.size() + " of the topic. The larger number will be used instead."); } } - this.topicsInUse.put(destination, partitions); - - DefaultKafkaProducerFactory producerFB = getProducerFactory(producerProperties); KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFB); if (this.producerListener != null) { kafkaTemplate.setProducerListener(this.producerListener); } - return new ProducerConfigurationMessageHandler(kafkaTemplate, destination, producerProperties, producerFB); - } - - @Override - protected String createProducerDestinationIfNecessary(String name, - ExtendedProducerProperties properties) { - if (this.logger.isInfoEnabled()) { - this.logger.info("Using kafka topic for outbound: " + name); - } - KafkaTopicUtils.validateTopicName(name); - createTopicsIfAutoCreateEnabledAndAdminUtilsPresent(name, properties.getPartitionCount()); - Collection partitions = getPartitionsForTopic(name, properties.getPartitionCount()); - if (properties.getPartitionCount() < partitions.size()) { - if (this.logger.isInfoEnabled()) { - this.logger.info("The `partitionCount` of the producer for topic " + name + " is " - + properties.getPartitionCount() + ", smaller than the actual partition count of " - + partitions.size() + " of the topic. The larger number will be used instead."); - } - } - this.topicsInUse.put(name, partitions); - return name; + return new ProducerConfigurationMessageHandler(kafkaTemplate, destination.getName(), producerProperties, producerFB); } private DefaultKafkaProducerFactory getProducerFactory( @@ -263,15 +192,28 @@ public class KafkaMessageChannelBinder extends } @Override - protected Collection createConsumerDestinationIfNecessary(String name, String group, - ExtendedConsumerProperties properties) { - KafkaTopicUtils.validateTopicName(name); - if (properties.getInstanceCount() == 0) { - throw new IllegalArgumentException("Instance count cannot be zero"); + @SuppressWarnings("unchecked") + protected MessageProducer createConsumerEndpoint(final ConsumerDestination destination, final String group, + ExtendedConsumerProperties properties) { + + boolean anonymous = !StringUtils.hasText(group); + Assert.isTrue(!anonymous || !properties.getExtension().isEnableDlq(), + "DLQ support is not available for anonymous subscriptions"); + String consumerGroup = anonymous ? "anonymous." + UUID.randomUUID().toString() : group; + Map props = getConsumerConfig(anonymous, consumerGroup); + if (!ObjectUtils.isEmpty(properties.getExtension().getConfiguration())) { + props.putAll(properties.getExtension().getConfiguration()); } + final ConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory<>(props); int partitionCount = properties.getInstanceCount() * properties.getConcurrency(); - createTopicsIfAutoCreateEnabledAndAdminUtilsPresent(name, partitionCount); - Collection allPartitions = getPartitionsForTopic(name, partitionCount); + + Collection allPartitions = provisioningProvider.getPartitionsForTopic(partitionCount, + new Callable>() { + @Override + public Collection call() throws Exception { + return consumerFactory.createConsumer().partitionsFor(destination.getName()); + } + }); Collection listenedPartitions; @@ -288,29 +230,13 @@ public class KafkaMessageChannelBinder extends } } } - this.topicsInUse.put(name, listenedPartitions); - return listenedPartitions; - } + this.topicsInUse.put(destination.getName(), listenedPartitions); - @Override - @SuppressWarnings("unchecked") - protected MessageProducer createConsumerEndpoint(String name, String group, Collection destination, - ExtendedConsumerProperties properties) { - boolean anonymous = !StringUtils.hasText(group); - Assert.isTrue(!anonymous || !properties.getExtension().isEnableDlq(), - "DLQ support is not available for anonymous subscriptions"); - String consumerGroup = anonymous ? "anonymous." + UUID.randomUUID().toString() : group; - Map props = getConsumerConfig(anonymous, consumerGroup); - if (!ObjectUtils.isEmpty(properties.getExtension().getConfiguration())) { - props.putAll(properties.getExtension().getConfiguration()); - } - ConsumerFactory consumerFactory = new DefaultKafkaConsumerFactory<>(props); - Collection listenedPartitions = destination; Assert.isTrue(!CollectionUtils.isEmpty(listenedPartitions), "A list of partitions must be provided"); final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets( listenedPartitions); final ContainerProperties containerProperties = - anonymous || properties.getExtension().isAutoRebalanceEnabled() ? new ContainerProperties(name) + anonymous || properties.getExtension().isAutoRebalanceEnabled() ? new ContainerProperties(destination.getName()) : new ContainerProperties(topicPartitionInitialOffsets); int concurrency = Math.min(properties.getConcurrency(), listenedPartitions.size()); final ConcurrentMessageListenerContainer messageListenerContainer = @@ -342,8 +268,8 @@ public class KafkaMessageChannelBinder extends final RetryTemplate retryTemplate = buildRetryTemplate(properties); kafkaMessageDrivenChannelAdapter.setRetryTemplate(retryTemplate); if (properties.getExtension().isEnableDlq()) { - final String dlqTopic = "error." + name + "." + group; - initDlqProducer(); + DefaultKafkaProducerFactory producerFactory = getProducerFactory(new ExtendedProducerProperties<>(new KafkaProducerProperties())); + final KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFactory); messageListenerContainer.getContainerProperties().setErrorHandler(new ErrorHandler() { @Override @@ -352,29 +278,30 @@ public class KafkaMessageChannelBinder extends : null; final byte[] payload = message.value() != null ? Utils.toArray(ByteBuffer.wrap((byte[]) message.value())) : null; - KafkaMessageChannelBinder.this.dlqProducer.send(new ProducerRecord<>(dlqTopic, key, payload), - new Callback() { + ListenableFuture> sentDlq = kafkaTemplate.send("error." + destination.getName() + "." + group, + message.partition(), key, payload); + sentDlq.addCallback(new ListenableFutureCallback>() { + StringBuilder sb = new StringBuilder().append(" a message with key='") + .append(toDisplayString(ObjectUtils.nullSafeToString(key), 50)).append("'") + .append(" and payload='") + .append(toDisplayString(ObjectUtils.nullSafeToString(payload), 50)) + .append("'").append(" received from ") + .append(message.partition()); - @Override - public void onCompletion(RecordMetadata metadata, Exception exception) { - StringBuffer messageLog = new StringBuffer(); - messageLog.append(" a message with key='" - + toDisplayString(ObjectUtils.nullSafeToString(key), 50) + "'"); - messageLog.append(" and payload='" - + toDisplayString(ObjectUtils.nullSafeToString(payload), 50) + "'"); - messageLog.append(" received from " + message.partition()); - if (exception != null) { - KafkaMessageChannelBinder.this.logger.error( - "Error sending to DLQ" + messageLog.toString(), exception); - } - else { - if (KafkaMessageChannelBinder.this.logger.isDebugEnabled()) { - KafkaMessageChannelBinder.this.logger.debug( - "Sent to DLQ " + messageLog.toString()); - } - } - } - }); + @Override + public void onFailure(Throwable ex) { + KafkaMessageChannelBinder.this.logger.error( + "Error sending to DLQ" + sb.toString(), ex); + } + + @Override + public void onSuccess(SendResult result) { + if (KafkaMessageChannelBinder.this.logger.isDebugEnabled()) { + KafkaMessageChannelBinder.this.logger.debug( + "Sent to DLQ " + sb.toString()); + } + } + }); } }); } @@ -416,146 +343,6 @@ public class KafkaMessageChannelBinder extends return topicPartitionInitialOffsets; } - private void createTopicsIfAutoCreateEnabledAndAdminUtilsPresent(final String topicName, final int partitionCount) { - if (this.configurationProperties.isAutoCreateTopics() && adminUtilsOperation != null) { - createTopicAndPartitions(topicName, partitionCount); - } - else if (this.configurationProperties.isAutoCreateTopics() && adminUtilsOperation == null) { - this.logger.warn("Auto creation of topics is enabled, but Kafka AdminUtils class is not present on the classpath. " + - "No topic will be created by the binder"); - } - else if (!this.configurationProperties.isAutoCreateTopics()) { - this.logger.info("Auto creation of topics is disabled."); - } - } - - /** - * Creates a Kafka topic if needed, or try to increase its partition count to the - * desired number. - */ - private void createTopicAndPartitions(final String topicName, final int partitionCount) { - - final ZkUtils zkUtils = ZkUtils.apply(this.configurationProperties.getZkConnectionString(), - this.configurationProperties.getZkSessionTimeout(), - this.configurationProperties.getZkConnectionTimeout(), - JaasUtils.isZkSecurityEnabled()); - try { - short errorCode = adminUtilsOperation.errorCodeFromTopicMetadata(topicName, zkUtils); - if (errorCode == ErrorMapping.NoError()) { - // only consider minPartitionCount for resizing if autoAddPartitions is true - int effectivePartitionCount = this.configurationProperties.isAutoAddPartitions() - ? Math.max(this.configurationProperties.getMinPartitionCount(), partitionCount) - : partitionCount; - int partitionSize = adminUtilsOperation.partitionSize(topicName, zkUtils); - - if (partitionSize < effectivePartitionCount) { - if (this.configurationProperties.isAutoAddPartitions()) { - adminUtilsOperation.invokeAddPartitions(zkUtils, topicName, effectivePartitionCount, null, false); - } - else { - throw new BinderException("The number of expected partitions was: " + partitionCount + ", but " - + partitionSize + (partitionSize > 1 ? " have " : " has ") + "been found instead." - + "Consider either increasing the partition count of the topic or enabling " + - "`autoAddPartitions`"); - } - } - } - else if (errorCode == ErrorMapping.UnknownTopicOrPartitionCode()) { - // always consider minPartitionCount for topic creation - final int effectivePartitionCount = Math.max(this.configurationProperties.getMinPartitionCount(), - partitionCount); - - this.metadataRetryOperations.execute(new RetryCallback() { - - @Override - public Object doWithRetry(RetryContext context) throws RuntimeException { - - 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; - } - }); - } - else { - throw new BinderException("Error fetching Kafka topic metadata: ", - ErrorMapping.exceptionFor(errorCode)); - } - } - finally { - zkUtils.close(); - } - } - - private Collection getPartitionsForTopic(final String topicName, final int partitionCount) { - try { - return this.metadataRetryOperations - .execute(new RetryCallback, Exception>() { - - @Override - public Collection doWithRetry(RetryContext context) throws Exception { - Collection partitions = - getProducerFactory( - new ExtendedProducerProperties<>(new KafkaProducerProperties())) - .createProducer().partitionsFor(topicName); - - // do a sanity check on the partition set - if (partitions.size() < partitionCount) { - throw new IllegalStateException("The number of expected partitions was: " - + partitionCount + ", but " + partitions.size() - + (partitions.size() > 1 ? " have " : " has ") + "been found instead"); - } - return partitions; - } - }); - } - catch (Exception e) { - this.logger.error("Cannot initialize Binder", e); - throw new BinderException("Cannot initialize binder:", e); - } - } - - private synchronized void initDlqProducer() { - try { - if (this.dlqProducer == null) { - synchronized (this) { - if (this.dlqProducer == null) { - // we can use the producer defaults as we do not need to tune - // performance - Map props = new HashMap<>(); - props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, - this.configurationProperties.getKafkaConnectionString()); - props.put(ProducerConfig.RETRIES_CONFIG, 0); - props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); - props.put(ProducerConfig.LINGER_MS_CONFIG, 1); - props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); - props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class); - props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class); - DefaultKafkaProducerFactory defaultKafkaProducerFactory = - new DefaultKafkaProducerFactory<>(props); - this.dlqProducer = defaultKafkaProducerFactory.createProducer(); - } - } - } - } - catch (Exception e) { - throw new RuntimeException("Cannot initialize DLQ producer:", e); - } - } - private String toDisplayString(String original, int maxCharacters) { if (original.length() <= maxCharacters) { return original; @@ -578,7 +365,7 @@ public class KafkaMessageChannelBinder extends setBeanFactory(KafkaMessageChannelBinder.this.getBeanFactory()); if (producerProperties.isPartitioned()) { SpelExpressionParser parser = new SpelExpressionParser(); - setPartitionIdExpression(parser.parseExpression("headers.partition")); + setPartitionIdExpression(parser.parseExpression("headers." + BinderHeaders.PARTITION_HEADER)); } if (producerProperties.getExtension().isSync()) { setSync(true); diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index 438cbeddd..fa7c25d4e 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -30,11 +30,14 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.kafka.KafkaBinderHealthIndicator; import org.springframework.cloud.stream.binder.kafka.KafkaBinderJaasInitializerListener; -import org.springframework.cloud.stream.binder.kafka.KafkaExtendedBindingProperties; import org.springframework.cloud.stream.binder.kafka.KafkaMessageChannelBinder; import org.springframework.cloud.stream.binder.kafka.admin.AdminUtilsOperation; import org.springframework.cloud.stream.binder.kafka.admin.Kafka09AdminUtilsOperation; import org.springframework.cloud.stream.binder.kafka.admin.Kafka10AdminUtilsOperation; +import org.springframework.cloud.stream.binder.kafka.properties.JaasLoginModuleConfiguration; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaExtendedBindingProperties; +import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.cloud.stream.config.codec.kryo.KryoCodecAutoConfiguration; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationListener; @@ -82,14 +85,18 @@ public class KafkaBinderConfiguration { @Autowired (required = false) private AdminUtilsOperation adminUtilsOperation; + @Bean + KafkaTopicProvisioner provisioningProvider() { + return new KafkaTopicProvisioner(this.configurationProperties, this.adminUtilsOperation); + } + @Bean KafkaMessageChannelBinder kafkaMessageChannelBinder() { KafkaMessageChannelBinder kafkaMessageChannelBinder = new KafkaMessageChannelBinder( - this.configurationProperties); + this.configurationProperties, provisioningProvider()); kafkaMessageChannelBinder.setCodec(this.codec); kafkaMessageChannelBinder.setProducerListener(producerListener); kafkaMessageChannelBinder.setExtendedBindingProperties(this.kafkaExtendedBindingProperties); - kafkaMessageChannelBinder.setAdminUtilsOperation(adminUtilsOperation); return kafkaMessageChannelBinder; } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java index e928575fd..043541dbf 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/AbstractKafkaTestBinder.java @@ -23,6 +23,8 @@ import com.esotericsoftware.kryo.Registration; import org.springframework.cloud.stream.binder.AbstractTestBinder; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; import org.springframework.integration.codec.Codec; import org.springframework.integration.codec.kryo.KryoRegistrar; import org.springframework.integration.codec.kryo.PojoCodec; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java index 9556faf43..f45c8b7a4 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java @@ -43,7 +43,9 @@ import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.Spy; import org.springframework.cloud.stream.binder.kafka.admin.Kafka10AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +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; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.kafka.core.ConsumerFactory; @@ -55,7 +57,6 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.support.MessageBuilder; -import org.springframework.retry.RetryOperations; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.fail; @@ -112,12 +113,6 @@ public class Kafka10BinderTests extends KafkaBinderTests { return consumerFactory().createConsumer().partitionsFor(topic).size(); } - @Override - @SuppressWarnings("unchecked") - protected void setMetadataRetryOperations(Binder binder, RetryOperations retryOperations) { - ((Kafka10TestBinder) binder).getBinder().setMetadataRetryOperations(retryOperations); - } - @Override protected ZkUtils getZkUtils(KafkaBinderConfigurationProperties kafkaBinderConfigurationProperties) { final ZkClient zkClient = new ZkClient(kafkaBinderConfigurationProperties.getZkConnectionString(), diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java index 9732bed6d..0f4212570 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java @@ -16,8 +16,10 @@ package org.springframework.cloud.stream.binder.kafka; +import org.springframework.cloud.stream.binder.kafka.admin.AdminUtilsOperation; import org.springframework.cloud.stream.binder.kafka.admin.Kafka10AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.provisioning.KafkaTopicProvisioner; import org.springframework.context.support.GenericApplicationContext; import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; @@ -34,14 +36,19 @@ public class Kafka10TestBinder extends AbstractKafkaTestBinder { public Kafka10TestBinder(KafkaBinderConfigurationProperties binderConfiguration) { try { - KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(binderConfiguration); + AdminUtilsOperation adminUtilsOperation = new Kafka10AdminUtilsOperation(); + KafkaTopicProvisioner provisioningProvider = + new KafkaTopicProvisioner(binderConfiguration, adminUtilsOperation); + provisioningProvider.afterPropertiesSet(); + + KafkaMessageChannelBinder binder = new KafkaMessageChannelBinder(binderConfiguration, provisioningProvider); + binder.setCodec(getCodec()); ProducerListener producerListener = new LoggingProducerListener(); binder.setProducerListener(producerListener); GenericApplicationContext context = new GenericApplicationContext(); context.refresh(); binder.setApplicationContext(context); - binder.setAdminUtilsOperation(new Kafka10AdminUtilsOperation()); binder.afterPropertiesSet(); this.setBinder(binder); } diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java index 400557a0d..c24757f6d 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderTests.java @@ -38,7 +38,6 @@ import org.junit.Test; import org.springframework.beans.DirectFieldAccessor; import org.springframework.cloud.stream.binder.Binder; -import org.springframework.cloud.stream.binder.BinderException; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.DefaultBinding; import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; @@ -46,8 +45,12 @@ import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.PartitionCapableBinderTests; import org.springframework.cloud.stream.binder.PartitionTestSupport; import org.springframework.cloud.stream.binder.TestUtils; -import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; +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; +import org.springframework.cloud.stream.binder.kafka.utils.KafkaTopicUtils; import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.cloud.stream.provisioning.ProvisioningException; import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.IntegrationMessageHeaderAccessor; import org.springframework.integration.channel.DirectChannel; @@ -66,7 +69,6 @@ import org.springframework.messaging.MessagingException; import org.springframework.messaging.SubscribableChannel; import org.springframework.messaging.support.GenericMessage; import org.springframework.messaging.support.MessageBuilder; -import org.springframework.retry.RetryOperations; import org.springframework.retry.backoff.FixedBackOffPolicy; import org.springframework.retry.policy.SimpleRetryPolicy; import org.springframework.retry.support.RetryTemplate; @@ -107,8 +109,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); String testTopicName = "nonexisting" + System.currentTimeMillis(); @@ -1013,7 +1012,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); String testTopicName = "createdByBroker-" + System.currentTimeMillis(); @@ -1069,7 +1067,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); Binding binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties); @@ -1106,7 +1103,7 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); // this consumer must consume from partition 2 @@ -1188,7 +1184,6 @@ public abstract class KafkaBinderTests extends PartitionCapableBinderTests consumerProperties = createConsumerProperties(); Binding binding = binder.bindConsumer(testTopicName, "test", output, consumerProperties);