From 7308bd4991b8483434f2faf16e6537f9a1d379f2 Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Wed, 7 Sep 2016 16:46:25 -0400 Subject: [PATCH] Polishing Removed optional flag --- pom.xml | 47 +++++++++++++++++ spring-cloud-starter-stream-kafka/pom.xml | 32 +----------- .../main/resources/META-INF/spring.provides | 2 +- .../pom.xml | 13 +++-- .../binder/kafka/Kafka10BinderTests.java | 2 +- .../binder/kafka/Kafka10TestBinder.java | 2 +- .../src/main/asciidoc/overview.adoc | 28 ++++++---- spring-cloud-stream-binder-kafka/pom.xml | 51 ++++--------------- .../kafka/KafkaBinderHealthIndicator.java | 2 +- .../kafka/KafkaMessageChannelBinder.java | 2 +- .../KafkaBinderConfiguration.java | 2 +- .../KafkaBinderConfigurationProperties.java | 2 +- .../main/resources/META-INF/spring.binders | 2 +- .../binder/kafka/Kafka09BinderTests.java | 2 +- .../binder/kafka/Kafka09TestBinder.java | 2 +- .../stream/binder/kafka/KafkaBinderTests.java | 2 +- 16 files changed, 93 insertions(+), 100 deletions(-) rename spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/{configuration => config}/KafkaBinderConfiguration.java (98%) rename spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/{configuration => config}/KafkaBinderConfigurationProperties.java (98%) diff --git a/pom.xml b/pom.xml index 3d1218dcd..f818adfa9 100644 --- a/pom.xml +++ b/pom.xml @@ -59,6 +59,53 @@ + + + + + org.apache.kafka + kafka_2.11 + ${kafka.version} + + + org.slf4j + slf4j-log4j12 + + + log4j + log4j + + + + + org.apache.kafka + kafka-clients + ${kafka.version} + + + org.springframework.kafka + spring-kafka + ${spring-kafka.version} + + + org.springframework.integration + spring-integration-kafka + ${spring-integration-kafka.version} + + + org.springframework.kafka + spring-kafka-test + test + ${spring-kafka.version} + + + org.apache.kafka + kafka_2.11 + test + ${kafka.version} + + + spring diff --git a/spring-cloud-starter-stream-kafka/pom.xml b/spring-cloud-starter-stream-kafka/pom.xml index ae85746a4..b1d327de7 100644 --- a/spring-cloud-starter-stream-kafka/pom.xml +++ b/spring-cloud-starter-stream-kafka/pom.xml @@ -15,42 +15,12 @@ ${basedir}/../.. - 0.9.0.1 - 1.0.3.RELEASE - 2.0.1.RELEASE - org.springframework.cloud spring-cloud-stream-binder-kafka - ${project.version} - - - org.springframework.kafka - spring-kafka - ${spring-kafka.version} - - - org.apache.kafka - kafka_2.11 - ${kafka.version} - - - org.slf4j - slf4j-log4j12 - - - - - org.apache.kafka - kafka-clients - ${kafka.version} - - - org.springframework.integration - spring-integration-kafka - ${spring-integration-kafka.version} + 1.1.0.BUILD-SNAPSHOT diff --git a/spring-cloud-starter-stream-kafka/src/main/resources/META-INF/spring.provides b/spring-cloud-starter-stream-kafka/src/main/resources/META-INF/spring.provides index d9ca8793e..cc7cb9cc2 100644 --- a/spring-cloud-starter-stream-kafka/src/main/resources/META-INF/spring.provides +++ b/spring-cloud-starter-stream-kafka/src/main/resources/META-INF/spring.provides @@ -1 +1 @@ -provides: spring-cloud-starter-stream-kafka-0.10 \ No newline at end of file +provides: spring-cloud-starter-stream-kafka \ No newline at end of file diff --git a/spring-cloud-stream-binder-kafka-0.10-test/pom.xml b/spring-cloud-stream-binder-kafka-0.10-test/pom.xml index 92aa97f58..495540370 100644 --- a/spring-cloud-stream-binder-kafka-0.10-test/pom.xml +++ b/spring-cloud-stream-binder-kafka-0.10-test/pom.xml @@ -15,6 +15,10 @@ ${basedir}/../.. + 0.10.0.0 1.1.0.M1 2.0.1.RELEASE @@ -24,19 +28,17 @@ org.springframework.cloud spring-cloud-stream-binder-kafka - ${project.version} + 1.1.0.BUILD-SNAPSHOT test org.springframework.kafka spring-kafka - ${spring-kafka.version} test org.apache.kafka kafka_2.11 - ${kafka.version} test @@ -48,24 +50,21 @@ org.apache.kafka kafka-clients - ${kafka.version} test org.springframework.kafka spring-kafka-test test - ${spring-kafka.version} org.springframework.integration spring-integration-kafka - ${spring-integration-kafka.version} org.springframework.cloud spring-cloud-stream-binder-kafka - ${project.version} + 1.1.0.BUILD-SNAPSHOT test-jar test diff --git a/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java b/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java index 26f3533bc..e6b6f8299 100644 --- a/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java +++ b/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10BinderTests.java @@ -36,7 +36,7 @@ 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.configuration.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.support.KafkaHeaders; diff --git a/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java b/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java index 82329db49..9732bed6d 100644 --- a/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java +++ b/spring-cloud-stream-binder-kafka-0.10-test/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka10TestBinder.java @@ -17,7 +17,7 @@ package org.springframework.cloud.stream.binder.kafka; import org.springframework.cloud.stream.binder.kafka.admin.Kafka10AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.configuration.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.context.support.GenericApplicationContext; import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; 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 ef6b8b823..53f4305dd 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 @@ -121,11 +121,12 @@ The following properties are available for Kafka consumers only and must be prefixed with `spring.cloud.stream.kafka.bindings..consumer.`. autoRebalanceEnabled:: - When this is enabled, topic partitions will be automatically rebalanced between various consumers by the broker. - If set to `false`, it will trigger partitions to be statically allocated by the binder. -When setting false, in order for static allocation of partitions to take place, it also needs both `spring.cloud.stream.instaceCount` and -'spring.cloud.stream.instanceIndex' properties set appropriately. The property `spring.cloud.stream.instaceCount` must be greater than 1 in this case. - +When `true`, topic partitions will be automatically rebalanced between the members of a consumer group. +When `false`, each consumer will be assigned a fixed set of partitions based on `spring.cloud.stream.instanceCount` and `spring.cloud.stream.instanceIndex`. +This requires both `spring.cloud.stream.instanceCount` and `spring.cloud.stream.instanceIndex` properties to be set appropriately on each launched instance. +The property `spring.cloud.stream.instanceCount` must typically be greater than 1 in this case. ++ +Default: `true`. autoCommitOffset:: Whether to autocommit offsets when a message has been processed. If set to `false`, an `Acknowledgment` header will be available in the message headers for late acknowledgment. @@ -234,10 +235,11 @@ Usually applications may use principals that do not have administrative rights i In secure environments, we strongly recommend creating topics and managing ACLs administratively using Kafka tooling. ==== -[NOTE] -==== -In addition to supporting 0.9 based clients, Kafka binder can also work with 0.10 libs. In order to support this, when you create the project that contains your application, include `spring-cloud-starter-stream-kafka` as you normally would do for 0.9 based applications. -Then add these dependencies at the top of the dependencies section in the pom.xml file. +==== Using the binder with Apache Kafka 0.10 + +The binder also supports connecting to Kafka 0.10 brokers. +In order to support this, when you create the project that contains your application, include `spring-cloud-starter-stream-kafka` as you normally would do for 0.9 based applications. +Then add these dependencies at the top of the `` section in the pom.xml file to override the Apache Kafka, Spring Kafka, and Spring Integration Kafka with 0.10-compatible versions as in the following example: [source,xml] ---- @@ -264,4 +266,10 @@ Then add these dependencies at the top of the dependencies section in the pom.xm ---- -==== \ No newline at end of file +==== + +[NOTE] +==== +The versions above are provided only for the sake of the example. +For best results, we recommend using the most recent 0.10-compatible versions of the projects. +==== diff --git a/spring-cloud-stream-binder-kafka/pom.xml b/spring-cloud-stream-binder-kafka/pom.xml index 4d943ae71..de8c433cd 100644 --- a/spring-cloud-stream-binder-kafka/pom.xml +++ b/spring-cloud-stream-binder-kafka/pom.xml @@ -13,12 +13,6 @@ 1.1.0.BUILD-SNAPSHOT - - 0.9.0.1 - 1.0.3.RELEASE - 2.0.1.RELEASE - - org.springframework.boot @@ -58,23 +52,19 @@ org.springframework.kafka spring-kafka ${spring-kafka.version} - true org.apache.kafka kafka_2.11 - true org.apache.kafka kafka-clients - true org.springframework.kafka spring-kafka-test test - ${spring-kafka.version} org.apache.kafka @@ -82,39 +72,18 @@ test test + + org.springframework.kafka + spring-kafka + ${spring-kafka.version} + + + org.springframework.integration + spring-integration-kafka + ${spring-integration-kafka.version} + - - - - org.apache.kafka - kafka_2.11 - ${kafka.version} - - - org.slf4j - slf4j-log4j12 - - - log4j - log4j - - - - - org.apache.kafka - kafka_2.11 - test - ${kafka.version} - - - org.apache.kafka - kafka-clients - ${kafka.version} - - - - 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 8e35d5fe3..d36552ec1 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.configuration.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; /** * Health indicator for Kafka. 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 548a07b7b..c1969d366 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 @@ -50,7 +50,7 @@ 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.configuration.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.context.Lifecycle; import org.springframework.expression.common.LiteralExpression; import org.springframework.expression.spel.standard.SpelExpressionParser; diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/configuration/KafkaBinderConfiguration.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java similarity index 98% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/configuration/KafkaBinderConfiguration.java rename to spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java index 58353bb23..b926412f3 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/configuration/KafkaBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfiguration.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka.configuration; +package org.springframework.cloud.stream.binder.kafka.config; import java.lang.reflect.Method; diff --git a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/configuration/KafkaBinderConfigurationProperties.java b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java similarity index 98% rename from spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/configuration/KafkaBinderConfigurationProperties.java rename to spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java index 35ae8cda6..6c9e1467f 100644 --- a/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/configuration/KafkaBinderConfigurationProperties.java +++ b/spring-cloud-stream-binder-kafka/src/main/java/org/springframework/cloud/stream/binder/kafka/config/KafkaBinderConfigurationProperties.java @@ -14,7 +14,7 @@ * limitations under the License. */ -package org.springframework.cloud.stream.binder.kafka.configuration; +package org.springframework.cloud.stream.binder.kafka.config; import java.util.HashMap; import java.util.Map; diff --git a/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring.binders b/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring.binders index c6c1dc579..063a7400f 100644 --- a/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring.binders +++ b/spring-cloud-stream-binder-kafka/src/main/resources/META-INF/spring.binders @@ -1,2 +1,2 @@ kafka:\ -org.springframework.cloud.stream.binder.kafka.configuration.KafkaBinderConfiguration +org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfiguration diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09BinderTests.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09BinderTests.java index 7a0f6dabd..25ff24e09 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09BinderTests.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09BinderTests.java @@ -36,7 +36,7 @@ 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.Kafka09AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.configuration.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.kafka.core.ConsumerFactory; import org.springframework.kafka.core.DefaultKafkaConsumerFactory; import org.springframework.kafka.support.KafkaHeaders; diff --git a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java index c0ca54214..31ca93e5b 100644 --- a/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java +++ b/spring-cloud-stream-binder-kafka/src/test/java/org/springframework/cloud/stream/binder/kafka/Kafka09TestBinder.java @@ -17,7 +17,7 @@ package org.springframework.cloud.stream.binder.kafka; import org.springframework.cloud.stream.binder.kafka.admin.Kafka09AdminUtilsOperation; -import org.springframework.cloud.stream.binder.kafka.configuration.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.context.support.GenericApplicationContext; import org.springframework.kafka.support.LoggingProducerListener; import org.springframework.kafka.support.ProducerListener; 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 599a155bc..eabf0bd0f 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 @@ -40,7 +40,7 @@ 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.configuration.KafkaBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.kafka.config.KafkaBinderConfigurationProperties; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.context.support.GenericApplicationContext; import org.springframework.integration.IntegrationMessageHeaderAccessor;