diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java index 6f0d26d2b..e3d5ab739 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaConsumerProperties.java @@ -19,6 +19,8 @@ package org.springframework.cloud.stream.binder.kafka.properties; import java.util.HashMap; import java.util.Map; +import org.springframework.cloud.stream.config.MergableProperties; + /** * @author Marius Bogoevici * @author Ilayaperumal Gopinathan @@ -29,7 +31,7 @@ import java.util.Map; * Thanks to Laszlo Szabo for providing the initial patch for generic property support. *

*/ -public class KafkaConsumerProperties { +public class KafkaConsumerProperties implements MergableProperties { public enum StartOffset { earliest(-2L), diff --git a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaExtendedBindingProperties.java b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaExtendedBindingProperties.java index cac6d675a..ba0273819 100644 --- a/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaExtendedBindingProperties.java +++ b/spring-cloud-stream-binder-kafka-core/src/main/java/org/springframework/cloud/stream/binder/kafka/properties/KafkaExtendedBindingProperties.java @@ -25,11 +25,14 @@ import org.springframework.cloud.stream.binder.ExtendedBindingProperties; /** * @author Marius Bogoevici * @author Gary Russell + * @author Soby Chacko */ @ConfigurationProperties("spring.cloud.stream.kafka") public class KafkaExtendedBindingProperties implements ExtendedBindingProperties { + private static final String DEFAULTS_PREFIX = "spring.cloud.stream.kafka.default"; + private Map bindings = new HashMap<>(); public Map getBindings() { @@ -82,4 +85,14 @@ public class KafkaExtendedBindingProperties } } + @Override + public String getDefaultsPrefix() { + return DEFAULTS_PREFIX; + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return KafkaBindingProperties.class; + } + } 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 ae35dd324..3874f18dc 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 @@ -21,6 +21,7 @@ import java.util.Map; import javax.validation.constraints.NotNull; +import org.springframework.cloud.stream.config.MergableProperties; import org.springframework.expression.Expression; /** @@ -28,7 +29,7 @@ import org.springframework.expression.Expression; * @author Henryk Konsek * @author Gary Russell */ -public class KafkaProducerProperties { +public class KafkaProducerProperties implements MergableProperties { private int bufferSize = 16384; diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java index 39f49419f..5d51e1bbe 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java @@ -90,4 +90,14 @@ public class GlobalKTableBinder extends throw new UnsupportedOperationException("No producer binding is allowed and therefore no properties"); } + @Override + public String getDefaultsPrefix() { + return this.kafkaStreamsExtendedBindingProperties.getDefaultsPrefix(); + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return this.kafkaStreamsExtendedBindingProperties.getExtendedPropertiesEntryClass(); + } + } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java index 99a5849e7..57e8b740d 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java @@ -132,4 +132,14 @@ class KStreamBinder extends public void setKafkaStreamsExtendedBindingProperties(KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties) { this.kafkaStreamsExtendedBindingProperties = kafkaStreamsExtendedBindingProperties; } + + @Override + public String getDefaultsPrefix() { + return this.kafkaStreamsExtendedBindingProperties.getDefaultsPrefix(); + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return this.kafkaStreamsExtendedBindingProperties.getExtendedPropertiesEntryClass(); + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java index 1f72846e5..5ccbf0224 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java @@ -88,4 +88,14 @@ class KTableBinder extends public KafkaStreamsProducerProperties getExtendedProducerProperties(String channelName) { return this.kafkaStreamsExtendedBindingProperties.getExtendedProducerProperties(channelName); } + + @Override + public String getDefaultsPrefix() { + return this.kafkaStreamsExtendedBindingProperties.getDefaultsPrefix(); + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return this.kafkaStreamsExtendedBindingProperties.getExtendedPropertiesEntryClass(); + } } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsExtendedBindingProperties.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsExtendedBindingProperties.java index d21f95eed..a4bc4f71f 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsExtendedBindingProperties.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/properties/KafkaStreamsExtendedBindingProperties.java @@ -58,4 +58,14 @@ public class KafkaStreamsExtendedBindingProperties return new KafkaStreamsProducerProperties(); } } + + @Override + public String getDefaultsPrefix() { + return null; + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return KafkaStreamsBindingProperties.class; + } } 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 f09ab324c..dc507259f 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 @@ -214,6 +214,16 @@ public class KafkaMessageChannelBinder extends return this.extendedBindingProperties.getExtendedProducerProperties(channelName); } + @Override + public String getDefaultsPrefix() { + return this.extendedBindingProperties.getDefaultsPrefix(); + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return this.extendedBindingProperties.getExtendedPropertiesEntryClass(); + } + @Override protected MessageHandler createProducerMessageHandler(final ProducerDestination destination, ExtendedProducerProperties producerProperties, MessageChannel errorChannel) diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java new file mode 100644 index 000000000..2c18b7c04 --- /dev/null +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java @@ -0,0 +1,155 @@ +/* + * Copyright 2018 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.cloud.stream.binder.kafka.integration; + +import org.junit.AfterClass; +import org.junit.BeforeClass; +import org.junit.ClassRule; +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.context.SpringBootTest; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.annotation.Input; +import org.springframework.cloud.stream.annotation.Output; +import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binder.Binder; +import org.springframework.cloud.stream.binder.BinderFactory; +import org.springframework.cloud.stream.binder.ConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; +import org.springframework.cloud.stream.binder.ProducerProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaConsumerProperties; +import org.springframework.cloud.stream.binder.kafka.properties.KafkaProducerProperties; +import org.springframework.cloud.stream.messaging.Processor; +import org.springframework.cloud.stream.messaging.Sink; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.kafka.test.rule.EmbeddedKafkaRule; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.SubscribableChannel; +import org.springframework.messaging.handler.annotation.SendTo; +import org.springframework.test.context.junit4.SpringRunner; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * @author Soby Chacko + */ +@RunWith(SpringRunner.class) +@SpringBootTest(webEnvironment = SpringBootTest.WebEnvironment.NONE, +properties = {"spring.cloud.stream.kafka.bindings.output.producer.configuration.key.serializer=FooSerializer.class", + "spring.cloud.stream.kafka.default.producer.configuration.key.serializer=BarSerializer.class", + "spring.cloud.stream.kafka.default.producer.configuration.value.serializer=BarSerializer.class", + "spring.cloud.stream.kafka.bindings.input.consumer.configuration.key.serializer=FooSerializer.class", + "spring.cloud.stream.kafka.default.consumer.configuration.key.serializer=BarSerializer.class", + "spring.cloud.stream.kafka.default.consumer.configuration.value.serializer=BarSerializer.class", + "spring.cloud.stream.kafka.default.producer.configuration.foo=bar", + "spring.cloud.stream.kafka.bindings.output.producer.configuration.foo=bindingSpecificPropertyShouldWinOverDefault"}) +public class KafkaBinderExtendedPropertiesTest { + + private static final String KAFKA_BROKERS_PROPERTY = "spring.cloud.stream.kafka.binder.brokers"; + private static final String ZK_NODE_PROPERTY = "spring.cloud.stream.kafka.binder.zkNodes"; + + @ClassRule + public static EmbeddedKafkaRule kafkaEmbedded = new EmbeddedKafkaRule(1, true); + + @BeforeClass + public static void setup() { + System.setProperty(KAFKA_BROKERS_PROPERTY, kafkaEmbedded.getEmbeddedKafka().getBrokersAsString()); + System.setProperty(ZK_NODE_PROPERTY, kafkaEmbedded.getEmbeddedKafka().getZookeeperConnectionString()); + } + + @AfterClass + public static void clean() { + System.clearProperty(KAFKA_BROKERS_PROPERTY); + System.clearProperty(ZK_NODE_PROPERTY); + } + + @Autowired + private ConfigurableApplicationContext context; + + @Test + public void testKafkaBinderExtendedProperties() { + + BinderFactory binderFactory = context.getBeanFactory().getBean(BinderFactory.class); + Binder kafkaBinder = + binderFactory.getBinder("kafka", MessageChannel.class); + + KafkaProducerProperties kafkaProducerProperties = + (KafkaProducerProperties)((ExtendedPropertiesBinder) kafkaBinder).getExtendedProducerProperties("output"); + + //binding "output" gets FooSerializer defined on the binding itself and BarSerializer through default property. + assertThat(kafkaProducerProperties.getConfiguration().get("key.serializer")).isEqualTo("FooSerializer.class"); + assertThat(kafkaProducerProperties.getConfiguration().get("value.serializer")).isEqualTo("BarSerializer.class"); + + assertThat(kafkaProducerProperties.getConfiguration().get("foo")).isEqualTo("bindingSpecificPropertyShouldWinOverDefault"); + + KafkaConsumerProperties kafkaConsumerProperties = + (KafkaConsumerProperties)((ExtendedPropertiesBinder) kafkaBinder).getExtendedConsumerProperties("input"); + + //binding "input" gets FooSerializer defined on the binding itself and BarSerializer through default property. + assertThat(kafkaConsumerProperties.getConfiguration().get("key.serializer")).isEqualTo("FooSerializer.class"); + assertThat(kafkaConsumerProperties.getConfiguration().get("value.serializer")).isEqualTo("BarSerializer.class"); + + KafkaProducerProperties customKafkaProducerProperties = + (KafkaProducerProperties)((ExtendedPropertiesBinder) kafkaBinder).getExtendedProducerProperties("custom-out"); + + //binding "output" gets BarSerializer and BarSerializer for ker.serializer/value.serializer through default properties. + assertThat(customKafkaProducerProperties.getConfiguration().get("key.serializer")).isEqualTo("BarSerializer.class"); + assertThat(customKafkaProducerProperties.getConfiguration().get("value.serializer")).isEqualTo("BarSerializer.class"); + + //through default properties. + assertThat(customKafkaProducerProperties.getConfiguration().get("foo")).isEqualTo("bar"); + + KafkaConsumerProperties customKafkaConsumerProperties = + (KafkaConsumerProperties)((ExtendedPropertiesBinder) kafkaBinder).getExtendedConsumerProperties("custom-in"); + + //binding "input" gets BarSerializer and BarSerializer for ker.serializer/value.serializer through default properties. + assertThat(customKafkaConsumerProperties.getConfiguration().get("key.serializer")).isEqualTo("BarSerializer.class"); + assertThat(customKafkaConsumerProperties.getConfiguration().get("value.serializer")).isEqualTo("BarSerializer.class"); + } + + @EnableBinding(CustomBindingForExtendedPropertyTesting.class) + @EnableAutoConfiguration + public static class KafkaMetricsTestConfig { + + @StreamListener(Sink.INPUT) + @SendTo(Processor.OUTPUT) + public String process(String payload) throws InterruptedException { + return payload; + } + + @StreamListener("custom-in") + @SendTo("custom-out") + public String processCustom(String payload) throws InterruptedException { + return payload; + } + + } + + interface CustomBindingForExtendedPropertyTesting extends Processor{ + + @Output("custom-out") + MessageChannel customOut(); + + @Input("custom-in") + SubscribableChannel customIn(); + + } + +}