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 5f0370a8f..df0c78ebb 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 @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 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. @@ -26,20 +26,8 @@ import javax.validation.constraints.NotNull; */ public class KafkaProducerProperties { - /** - * bufferSize property is deprecated. - * It is recommended to set compressionType as one of the per binding Kafka producer `configuration` properties. - * If using KafkaAutoConfiguration from Spring Boot 1.5.x, `spring.kafka.producer.batchSize` property can also be used. - */ - @Deprecated private int bufferSize = 16384; - /** - * compressionType property is deprecated. - * It is recommended to set compressionType as one of the per binding Kafka producer `configuration` properties. - * If using KafkaAutoConfiguration from Spring Boot 1.5.x, `spring.kafka.producer.compressionType` property can also be used. - */ - @Deprecated private CompressionType compressionType = CompressionType.none; private boolean sync; @@ -89,14 +77,9 @@ public class KafkaProducerProperties { this.configuration = configuration; } - @Deprecated - /** - * @see compressionType - */ public enum CompressionType { none, gzip, - snappy, - lz4 + snappy } } diff --git a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc index b628c66e8..89cc90fa6 100644 --- a/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc +++ b/spring-cloud-stream-binder-kafka-docs/src/main/asciidoc/overview.adoc @@ -41,10 +41,6 @@ Partitioning also maps directly to Apache Kafka partitions as well. This section contains the configuration options used by the Apache Kafka binder. -When using Spring Boot 1.5.x and above, one can user `KafkaProperties` (prefixed with 'spring.kafka') from `KafkaAutoConfiguration` to configure Kafka common, producer, consumer properties. - -Note: Any Spring Cloud Stream Kafka binder properties or the per binding Kafka producer/consumer properties get the precedence over the Spring Boot KafkaProperties. - For common configuration options and properties pertaining to binder, refer to the https://github.com/spring-cloud/spring-cloud-stream/blob/master/spring-cloud-stream-docs/src/main/asciidoc/spring-cloud-stream-overview.adoc#configuration-options[core docs]. === Kafka Binder Properties diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java index 3d660ddbc..6f59c5f0e 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderEnvironmentPostProcessor.java @@ -19,9 +19,6 @@ package org.springframework.cloud.stream.binder.kafka; import java.util.HashMap; import java.util.Map; -import org.apache.kafka.common.serialization.ByteArrayDeserializer; -import org.apache.kafka.common.serialization.ByteArraySerializer; - import org.springframework.boot.SpringApplication; import org.springframework.boot.env.EnvironmentPostProcessor; import org.springframework.core.env.ConfigurableEnvironment; @@ -29,67 +26,19 @@ import org.springframework.core.env.MapPropertySource; /** * An {@link EnvironmentPostProcessor} that sets some common configuration properties (log config etc.,) for Kafka - * binder. + * binder. * * @author Ilayaperumal Gopinathan */ public class KafkaBinderEnvironmentPostProcessor implements EnvironmentPostProcessor { - public final static String SPRING_KAFKA = "spring.kafka"; - - public final static String SPRING_KAFKA_PRODUCER = SPRING_KAFKA + ".producer"; - - public final static String SPRING_KAFKA_CONSUMER = SPRING_KAFKA + ".consumer"; - - public final static String SPRING_KAFKA_PRODUCER_KEY_SERIALIZER = SPRING_KAFKA_PRODUCER + "." + "keySerializer"; - - public final static String SPRING_KAFKA_PRODUCER_VALUE_SERIALIZER = SPRING_KAFKA_PRODUCER + "." + "valueSerializer"; - - public final static String SPRING_KAFKA_CONSUMER_KEY_DESERIALIZER = SPRING_KAFKA_CONSUMER + "." + "keyDeserializer"; - - public final static String SPRING_KAFKA_CONSUMER_VALUE_DESERIALIZER = SPRING_KAFKA_CONSUMER + "." + "valueDeserializer"; - - public final static String SPRING_KAFKA_BOOTSTRAP_SERVERS = SPRING_KAFKA + "." + "bootstrapServers"; - @Override public void postProcessEnvironment(ConfigurableEnvironment environment, SpringApplication application) { - Map logProperties = new HashMap<>(); - logProperties.put("logging.pattern.console", "%d{ISO8601} %5p %t %c{2}:%L - %m%n"); - logProperties.put("logging.level.org.I0Itec.zkclient", "ERROR"); - logProperties.put("logging.level.kafka.server.KafkaConfig", "ERROR"); - logProperties.put("logging.level.kafka.admin.AdminClient.AdminConfig", "ERROR"); - environment.getPropertySources().addLast(new MapPropertySource("kafkaBinderLogConfig", logProperties)); - Map binderConfig = new HashMap<>(); - if (environment.getProperty(SPRING_KAFKA_PRODUCER_KEY_SERIALIZER) != null) { - binderConfig.put(SPRING_KAFKA_PRODUCER_KEY_SERIALIZER, environment.getProperty(SPRING_KAFKA_PRODUCER_KEY_SERIALIZER)); - } - else { - binderConfig.put(SPRING_KAFKA_PRODUCER_KEY_SERIALIZER, ByteArraySerializer.class); - } - if (environment.getProperty(SPRING_KAFKA_PRODUCER_VALUE_SERIALIZER) != null) { - binderConfig.put(SPRING_KAFKA_PRODUCER_VALUE_SERIALIZER, environment.getProperty(SPRING_KAFKA_PRODUCER_VALUE_SERIALIZER)); - } - else { - binderConfig.put(SPRING_KAFKA_PRODUCER_VALUE_SERIALIZER, ByteArraySerializer.class); - } - if (environment.getProperty(SPRING_KAFKA_CONSUMER_KEY_DESERIALIZER) != null) { - binderConfig.put(SPRING_KAFKA_CONSUMER_KEY_DESERIALIZER, environment.getProperty(SPRING_KAFKA_CONSUMER_KEY_DESERIALIZER)); - } - else { - binderConfig.put(SPRING_KAFKA_CONSUMER_KEY_DESERIALIZER, ByteArrayDeserializer.class); - } - if (environment.getProperty(SPRING_KAFKA_CONSUMER_VALUE_DESERIALIZER) != null) { - binderConfig.put(SPRING_KAFKA_CONSUMER_VALUE_DESERIALIZER, environment.getProperty(SPRING_KAFKA_CONSUMER_VALUE_DESERIALIZER)); - } - else { - binderConfig.put(SPRING_KAFKA_CONSUMER_VALUE_DESERIALIZER, ByteArrayDeserializer.class); - } - if (environment.getProperty(SPRING_KAFKA_BOOTSTRAP_SERVERS) != null) { - binderConfig.put(SPRING_KAFKA_BOOTSTRAP_SERVERS, environment.getProperty(SPRING_KAFKA_BOOTSTRAP_SERVERS)); - } - else { - binderConfig.put(SPRING_KAFKA_BOOTSTRAP_SERVERS, ""); - } - environment.getPropertySources().addLast(new MapPropertySource("kafkaBinderConfig", binderConfig)); + Map propertiesToAdd = new HashMap<>(); + propertiesToAdd.put("logging.pattern.console", "%d{ISO8601} %5p %t %c{2}:%L - %m%n"); + propertiesToAdd.put("logging.level.org.I0Itec.zkclient", "ERROR"); + propertiesToAdd.put("logging.level.kafka.server.KafkaConfig", "ERROR"); + propertiesToAdd.put("logging.level.kafka.admin.AdminClient.AdminConfig", "ERROR"); + environment.getPropertySources().addLast(new MapPropertySource("kafkaBinderLogConfig", propertiesToAdd)); } } 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 381dfcd98..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 @@ -33,7 +33,6 @@ import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.apache.kafka.common.utils.Utils; -import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.BinderHeaders; @@ -98,10 +97,8 @@ public class KafkaMessageChannelBinder extends private final Map> topicsInUse = new HashMap<>(); - private KafkaProperties kafkaProperties; - public KafkaMessageChannelBinder(KafkaBinderConfigurationProperties configurationProperties, - KafkaTopicProvisioner provisioningProvider) { + KafkaTopicProvisioner provisioningProvider) { super(false, headersToMap(configurationProperties), provisioningProvider); this.configurationProperties = configurationProperties; } @@ -130,10 +127,6 @@ public class KafkaMessageChannelBinder extends this.producerListener = producerListener; } - public void setKafkaProperties(KafkaProperties kafkaProperties) { - this.kafkaProperties = kafkaProperties; - } - Map> getTopicsInUse() { return this.topicsInUse; } @@ -150,7 +143,7 @@ public class KafkaMessageChannelBinder extends @Override protected MessageHandler createProducerMessageHandler(final ProducerDestination destination, - ExtendedProducerProperties producerProperties) throws Exception { + ExtendedProducerProperties producerProperties) throws Exception { final DefaultKafkaProducerFactory producerFB = getProducerFactory(producerProperties); Collection partitions = provisioningProvider.getPartitionsForTopic(producerProperties.getPartitionCount(), new Callable>() { @@ -178,27 +171,20 @@ public class KafkaMessageChannelBinder extends private DefaultKafkaProducerFactory getProducerFactory( ExtendedProducerProperties producerProperties) { Map props = new HashMap<>(); + if (!ObjectUtils.isEmpty(configurationProperties.getConfiguration())) { + props.putAll(configurationProperties.getConfiguration()); + } + props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.configurationProperties.getKafkaConnectionString()); props.put(ProducerConfig.RETRIES_CONFIG, 0); + props.put(ProducerConfig.BATCH_SIZE_CONFIG, String.valueOf(producerProperties.getExtension().getBufferSize())); 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); - props.put(ProducerConfig.BATCH_SIZE_CONFIG, String.valueOf(producerProperties.getExtension().getBufferSize())); props.put(ProducerConfig.ACKS_CONFIG, String.valueOf(this.configurationProperties.getRequiredAcks())); props.put(ProducerConfig.LINGER_MS_CONFIG, String.valueOf(producerProperties.getExtension().getBatchTimeout())); props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, producerProperties.getExtension().getCompressionType().toString()); - props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.configurationProperties.getKafkaConnectionString()); - if (this.kafkaProperties != null) { - if (!this.kafkaProperties.getBootstrapServers().isEmpty()) { - props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.kafkaProperties.getBootstrapServers()); - } - props.putAll(this.kafkaProperties.getProducer().buildProperties()); - props.putAll(this.kafkaProperties.getProperties()); - } - if (!ObjectUtils.isEmpty(configurationProperties.getConfiguration())) { - props.putAll(configurationProperties.getConfiguration()); - } if (!ObjectUtils.isEmpty(producerProperties.getExtension().getConfiguration())) { props.putAll(producerProperties.getExtension().getConfiguration()); } @@ -208,13 +194,19 @@ public class KafkaMessageChannelBinder extends @Override @SuppressWarnings("unchecked") protected MessageProducer createConsumerEndpoint(final ConsumerDestination destination, final String group, - ExtendedConsumerProperties extendedConsumerProperties) { + ExtendedConsumerProperties properties) { + boolean anonymous = !StringUtils.hasText(group); - Assert.isTrue(!anonymous || !extendedConsumerProperties.getExtension().isEnableDlq(), + Assert.isTrue(!anonymous || !properties.getExtension().isEnableDlq(), "DLQ support is not available for anonymous subscriptions"); String consumerGroup = anonymous ? "anonymous." + UUID.randomUUID().toString() : group; - final ConsumerFactory consumerFactory = createKafkaConsumerFactory(anonymous, consumerGroup, extendedConsumerProperties); - int partitionCount = extendedConsumerProperties.getInstanceCount() * extendedConsumerProperties.getConcurrency(); + 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(); + Collection allPartitions = provisioningProvider.getPartitionsForTopic(partitionCount, new Callable>() { @Override @@ -222,28 +214,31 @@ public class KafkaMessageChannelBinder extends return consumerFactory.createConsumer().partitionsFor(destination.getName()); } }); + Collection listenedPartitions; - if (extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() || - extendedConsumerProperties.getInstanceCount() == 1) { + + if (properties.getExtension().isAutoRebalanceEnabled() || + properties.getInstanceCount() == 1) { listenedPartitions = allPartitions; } else { listenedPartitions = new ArrayList<>(); for (PartitionInfo partition : allPartitions) { // divide partitions across modules - if ((partition.partition() % extendedConsumerProperties.getInstanceCount()) == extendedConsumerProperties.getInstanceIndex()) { + if ((partition.partition() % properties.getInstanceCount()) == properties.getInstanceIndex()) { listenedPartitions.add(partition); } } } this.topicsInUse.put(destination.getName(), listenedPartitions); + Assert.isTrue(!CollectionUtils.isEmpty(listenedPartitions), "A list of partitions must be provided"); final TopicPartitionInitialOffset[] topicPartitionInitialOffsets = getTopicPartitionInitialOffsets( listenedPartitions); final ContainerProperties containerProperties = - anonymous || extendedConsumerProperties.getExtension().isAutoRebalanceEnabled() ? new ContainerProperties(destination.getName()) + anonymous || properties.getExtension().isAutoRebalanceEnabled() ? new ContainerProperties(destination.getName()) : new ContainerProperties(topicPartitionInitialOffsets); - int concurrency = Math.min(extendedConsumerProperties.getConcurrency(), listenedPartitions.size()); + int concurrency = Math.min(properties.getConcurrency(), listenedPartitions.size()); final ConcurrentMessageListenerContainer messageListenerContainer = new ConcurrentMessageListenerContainer( consumerFactory, containerProperties) { @@ -254,8 +249,8 @@ public class KafkaMessageChannelBinder extends } }; messageListenerContainer.setConcurrency(concurrency); - messageListenerContainer.getContainerProperties().setAckOnError(isAutoCommitOnError(extendedConsumerProperties)); - if (!extendedConsumerProperties.getExtension().isAutoCommitOffset()) { + messageListenerContainer.getContainerProperties().setAckOnError(isAutoCommitOnError(properties)); + if (!properties.getExtension().isAutoCommitOffset()) { messageListenerContainer.getContainerProperties().setAckMode(AbstractMessageListenerContainer.AckMode.MANUAL); } if (this.logger.isDebugEnabled()) { @@ -270,9 +265,9 @@ public class KafkaMessageChannelBinder extends new KafkaMessageDrivenChannelAdapter<>( messageListenerContainer); kafkaMessageDrivenChannelAdapter.setBeanFactory(this.getBeanFactory()); - final RetryTemplate retryTemplate = buildRetryTemplate(extendedConsumerProperties); + final RetryTemplate retryTemplate = buildRetryTemplate(properties); kafkaMessageDrivenChannelAdapter.setRetryTemplate(retryTemplate); - if (extendedConsumerProperties.getExtension().isEnableDlq()) { + if (properties.getExtension().isEnableDlq()) { DefaultKafkaProducerFactory producerFactory = getProducerFactory(new ExtendedProducerProperties<>(new KafkaProducerProperties())); final KafkaTemplate kafkaTemplate = new KafkaTemplate<>(producerFactory); messageListenerContainer.getContainerProperties().setErrorHandler(new ErrorHandler() { @@ -313,30 +308,20 @@ public class KafkaMessageChannelBinder extends return kafkaMessageDrivenChannelAdapter; } - private ConsumerFactory createKafkaConsumerFactory(boolean anonymous, String consumerGroup, - ExtendedConsumerProperties consumerProperties) { + private Map getConsumerConfig(boolean anonymous, String consumerGroup) { Map props = new HashMap<>(); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class); - props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); - props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup); - props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, anonymous ? "latest" : "earliest"); - props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 100); - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, this.configurationProperties.getKafkaConnectionString()); - if (this.kafkaProperties != null) { - if (!this.kafkaProperties.getBootstrapServers().isEmpty()) { - props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, this.kafkaProperties.getBootstrapServers()); - } - props.putAll(this.kafkaProperties.getConsumer().buildProperties()); - props.putAll(this.kafkaProperties.getProperties()); - } if (!ObjectUtils.isEmpty(configurationProperties.getConfiguration())) { props.putAll(configurationProperties.getConfiguration()); } - if (!ObjectUtils.isEmpty(consumerProperties.getExtension().getConfiguration())) { - props.putAll(consumerProperties.getExtension().getConfiguration()); - } - return new DefaultKafkaConsumerFactory<>(props); + props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, this.configurationProperties.getKafkaConnectionString()); + props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); + props.put(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup); + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, + anonymous ? "latest" : "earliest"); + props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 100); + return props; } private boolean isAutoCommitOnError(ExtendedConsumerProperties properties) { @@ -373,8 +358,8 @@ public class KafkaMessageChannelBinder extends private final DefaultKafkaProducerFactory producerFactory; private ProducerConfigurationMessageHandler(KafkaTemplate kafkaTemplate, String topic, - ExtendedProducerProperties producerProperties, - DefaultKafkaProducerFactory producerFactory) { + ExtendedProducerProperties producerProperties, + DefaultKafkaProducerFactory producerFactory) { super(kafkaTemplate); setTopicExpression(new LiteralExpression(topic)); setBeanFactory(KafkaMessageChannelBinder.this.getBeanFactory()); 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 4e0138951..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 @@ -26,7 +26,6 @@ import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.PropertyPlaceholderAutoConfiguration; import org.springframework.boot.autoconfigure.condition.ConditionalOnClass; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; -import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.binder.kafka.KafkaBinderHealthIndicator; @@ -40,6 +39,7 @@ import org.springframework.cloud.stream.binder.kafka.properties.KafkaBinderConfi 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; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Condition; @@ -79,8 +79,8 @@ public class KafkaBinderConfiguration { @Autowired private ProducerListener producerListener; - @Autowired(required = false) - private KafkaProperties kafkaProperties; + @Autowired + private ApplicationContext context; @Autowired (required = false) private AdminUtilsOperation adminUtilsOperation; @@ -97,7 +97,6 @@ public class KafkaBinderConfiguration { kafkaMessageChannelBinder.setCodec(this.codec); kafkaMessageChannelBinder.setProducerListener(producerListener); kafkaMessageChannelBinder.setExtendedBindingProperties(this.kafkaExtendedBindingProperties); - kafkaMessageChannelBinder.setKafkaProperties(kafkaProperties); return kafkaMessageChannelBinder; } @@ -140,7 +139,7 @@ public class KafkaBinderConfiguration { return AppInfoParser.getVersion().startsWith("0.10"); } } - + static class Kafka09Present implements Condition { @Override diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java index 7897c3123..2922ce6ef 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/KafkaBinderConfigurationTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2017 the original author or authors. + * Copyright 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. @@ -15,42 +15,25 @@ */ package org.springframework.cloud.stream.binder.kafka; -import java.lang.reflect.Field; -import java.lang.reflect.Method; -import java.util.ArrayList; -import java.util.List; -import java.util.Map; +import static org.junit.Assert.assertNotNull; + +import java.lang.reflect.Field; -import org.apache.kafka.common.serialization.LongDeserializer; -import org.apache.kafka.common.serialization.LongSerializer; import org.junit.Test; import org.junit.runner.RunWith; import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.autoconfigure.kafka.KafkaProperties; import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; -import org.springframework.cloud.stream.binder.ExtendedProducerProperties; import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration; -import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; -import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; -import org.springframework.context.annotation.Bean; -import org.springframework.kafka.core.DefaultKafkaConsumerFactory; -import org.springframework.kafka.core.DefaultKafkaProducerFactory; import org.springframework.kafka.support.ProducerListener; -import org.springframework.test.context.TestPropertySource; import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import org.springframework.util.ReflectionUtils; -import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertTrue; - /** * @author Ilayaperumal Gopinathan */ @RunWith(SpringJUnit4ClassRunner.class) -@SpringBootTest(classes = {KafkaBinderConfigurationTest.KafkaBinderConfigProperties.class, KafkaBinderConfiguration.class}) -@TestPropertySource(locations = "classpath:binder-config.properties") +@SpringBootTest(classes = KafkaBinderConfiguration.class) public class KafkaBinderConfigurationTest { @Autowired @@ -67,46 +50,4 @@ public class KafkaBinderConfigurationTest { producerListenerField, this.kafkaMessageChannelBinder); assertNotNull(producerListener); } - - @Test - public void testKafkaBinderConfiguration() throws Exception { - assertNotNull(this.kafkaMessageChannelBinder); - Field kafkaPropertiesField = ReflectionUtils.findField(KafkaMessageChannelBinder.class, "kafkaProperties", KafkaProperties.class); - ReflectionUtils.makeAccessible(kafkaPropertiesField); - KafkaProperties kafkaProperties = (KafkaProperties) ReflectionUtils.getField(kafkaPropertiesField, this.kafkaMessageChannelBinder); - assertNotNull(kafkaProperties); - ExtendedProducerProperties producerProperties = new ExtendedProducerProperties<>(new KafkaProducerProperties()); - Method getProducerFactoryMethod = KafkaMessageChannelBinder.class.getDeclaredMethod("getProducerFactory", ExtendedProducerProperties.class); - getProducerFactoryMethod.setAccessible(true); - DefaultKafkaProducerFactory producerFactory = (DefaultKafkaProducerFactory) getProducerFactoryMethod.invoke(this.kafkaMessageChannelBinder, producerProperties); - Field producerFactoryConfigField = ReflectionUtils.findField(DefaultKafkaProducerFactory.class, "configs", Map.class); - ReflectionUtils.makeAccessible(producerFactoryConfigField); - Map producerConfigs = (Map) ReflectionUtils.getField(producerFactoryConfigField, producerFactory); - assertTrue(producerConfigs.get("batch.size").equals(10)); - assertTrue(producerConfigs.get("key.serializer").equals(LongSerializer.class)); - assertTrue(producerConfigs.get("value.serializer").equals(LongSerializer.class)); - assertTrue(producerConfigs.get("compression.type").equals("snappy")); - List bootstrapServers = new ArrayList<>(); - bootstrapServers.add("10.98.09.199:9092"); - bootstrapServers.add("10.98.09.196:9092"); - assertTrue((((List) producerConfigs.get("bootstrap.servers")).containsAll(bootstrapServers))); - Method createKafkaConsumerFactoryMethod = KafkaMessageChannelBinder.class.getDeclaredMethod("createKafkaConsumerFactory", boolean.class, String.class, ExtendedConsumerProperties.class); - createKafkaConsumerFactoryMethod.setAccessible(true); - ExtendedConsumerProperties consumerProperties = new ExtendedConsumerProperties<>(new KafkaConsumerProperties()); - DefaultKafkaConsumerFactory consumerFactory = (DefaultKafkaConsumerFactory) createKafkaConsumerFactoryMethod.invoke(this.kafkaMessageChannelBinder, true, "test", consumerProperties); - Field consumerFactoryConfigField = ReflectionUtils.findField(DefaultKafkaConsumerFactory.class, "configs", Map.class); - ReflectionUtils.makeAccessible(consumerFactoryConfigField); - Map consumerConfigs = (Map) ReflectionUtils.getField(consumerFactoryConfigField, consumerFactory); - assertTrue(consumerConfigs.get("key.deserializer").equals(LongDeserializer.class)); - assertTrue(consumerConfigs.get("value.deserializer").equals(LongDeserializer.class)); - assertTrue((((List) consumerConfigs.get("bootstrap.servers")).containsAll(bootstrapServers))); - } - - public static class KafkaBinderConfigProperties { - - @Bean - KafkaProperties kafkaProperties() { - return new KafkaProperties(); - } - } } diff --git a/spring-cloud-stream-binder-kafka/src/test/resources/binder-config.properties b/spring-cloud-stream-binder-kafka/src/test/resources/binder-config.properties deleted file mode 100644 index 8ac39c0e5..000000000 --- a/spring-cloud-stream-binder-kafka/src/test/resources/binder-config.properties +++ /dev/null @@ -1,7 +0,0 @@ -spring.kafka.producer.keySerializer=org.apache.kafka.common.serialization.LongSerializer -spring.kafka.producer.valueSerializer=org.apache.kafka.common.serialization.LongSerializer -spring.kafka.consumer.keyDeserializer=org.apache.kafka.common.serialization.LongDeserializer -spring.kafka.consumer.valueDeserializer=org.apache.kafka.common.serialization.LongDeserializer -spring.kafka.producer.batchSize=10 -spring.kafka.bootstrapServers=10.98.09.199:9092,10.98.09.196:9092 -spring.kafka.producer.compressionType=snappy