From cc61d916b84d745fadda8feb57d89e8f4475bc5d Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Mon, 17 Sep 2018 20:09:18 -0400 Subject: [PATCH] Extended default properties Allow applications to configure default values for extended prducer and consumer properties across multiple bindings in order to avoid repetition. Address the changes in both Kafka and Kafka Streams binders. Add integration test to verify both binding specific and default extended properties. Resolves #444 Requires https://github.com/spring-cloud/spring-cloud-stream/pull/1477 --- .../properties/KafkaConsumerProperties.java | 4 +- .../KafkaExtendedBindingProperties.java | 13 ++ .../properties/KafkaProducerProperties.java | 3 +- .../kafka/streams/GlobalKTableBinder.java | 10 ++ .../binder/kafka/streams/KStreamBinder.java | 10 ++ .../binder/kafka/streams/KTableBinder.java | 10 ++ ...KafkaStreamsExtendedBindingProperties.java | 10 ++ .../kafka/KafkaMessageChannelBinder.java | 10 ++ .../KafkaBinderExtendedPropertiesTest.java | 155 ++++++++++++++++++ 9 files changed, 223 insertions(+), 2 deletions(-) create mode 100644 spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/integration/KafkaBinderExtendedPropertiesTest.java 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(); + + } + +}