diff --git a/kafka-streams/kafka-streams-branching-sample/.gitignore b/kafka-streams/kafka-streams-branching/.gitignore similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/.gitignore rename to kafka-streams/kafka-streams-branching/.gitignore diff --git a/kafka-streams/kafka-streams-branching-sample/.mvn/wrapper/maven-wrapper.jar b/kafka-streams/kafka-streams-branching/.mvn/wrapper/maven-wrapper.jar similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/.mvn/wrapper/maven-wrapper.jar rename to kafka-streams/kafka-streams-branching/.mvn/wrapper/maven-wrapper.jar diff --git a/kafka-streams/kafka-streams-branching-sample/.mvn/wrapper/maven-wrapper.properties b/kafka-streams/kafka-streams-branching/.mvn/wrapper/maven-wrapper.properties similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/.mvn/wrapper/maven-wrapper.properties rename to kafka-streams/kafka-streams-branching/.mvn/wrapper/maven-wrapper.properties diff --git a/kafka-streams/kafka-streams-branching-sample/README.adoc b/kafka-streams/kafka-streams-branching/README.adoc similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/README.adoc rename to kafka-streams/kafka-streams-branching/README.adoc diff --git a/kafka-streams/kafka-streams-branching-sample/docker/docker-compose.yml b/kafka-streams/kafka-streams-branching/docker/docker-compose.yml similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/docker/docker-compose.yml rename to kafka-streams/kafka-streams-branching/docker/docker-compose.yml diff --git a/kafka-streams/kafka-streams-branching-sample/docker/start-kafka-shell.sh b/kafka-streams/kafka-streams-branching/docker/start-kafka-shell.sh similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/docker/start-kafka-shell.sh rename to kafka-streams/kafka-streams-branching/docker/start-kafka-shell.sh diff --git a/kafka-streams/kafka-streams-branching-sample/mvnw b/kafka-streams/kafka-streams-branching/mvnw similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/mvnw rename to kafka-streams/kafka-streams-branching/mvnw diff --git a/kafka-streams/kafka-streams-branching-sample/mvnw.cmd b/kafka-streams/kafka-streams-branching/mvnw.cmd similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/mvnw.cmd rename to kafka-streams/kafka-streams-branching/mvnw.cmd diff --git a/kafka-streams/kafka-streams-branching-sample/pom.xml b/kafka-streams/kafka-streams-branching/pom.xml similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/pom.xml rename to kafka-streams/kafka-streams-branching/pom.xml diff --git a/kafka-streams/kafka-streams-branching-sample/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java b/kafka-streams/kafka-streams-branching/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java rename to kafka-streams/kafka-streams-branching/src/main/java/kafka/streams/branching/KafkaStreamsBranchingSample.java diff --git a/kafka-streams/kafka-streams-branching-sample/src/main/resources/application.yml b/kafka-streams/kafka-streams-branching/src/main/resources/application.yml similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/src/main/resources/application.yml rename to kafka-streams/kafka-streams-branching/src/main/resources/application.yml diff --git a/kafka-streams/kafka-streams-branching-sample/src/main/resources/logback.xml b/kafka-streams/kafka-streams-branching/src/main/resources/logback.xml similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/src/main/resources/logback.xml rename to kafka-streams/kafka-streams-branching/src/main/resources/logback.xml diff --git a/kafka-streams/kafka-streams-branching-sample/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java b/kafka-streams/kafka-streams-branching/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java similarity index 100% rename from kafka-streams/kafka-streams-branching-sample/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java rename to kafka-streams/kafka-streams-branching/src/test/java/kafka/streams/branching/KafkaStreamsBranchingSampleTests.java diff --git a/kafka-streams/kafka-streams-interactive-query/.gitignore b/kafka-streams/kafka-streams-interactive-query-advanced/.gitignore similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/.gitignore rename to kafka-streams/kafka-streams-interactive-query-advanced/.gitignore diff --git a/kafka-streams/kafka-streams-interactive-query/.mvn/wrapper/maven-wrapper.jar b/kafka-streams/kafka-streams-interactive-query-advanced/.mvn/wrapper/maven-wrapper.jar similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/.mvn/wrapper/maven-wrapper.jar rename to kafka-streams/kafka-streams-interactive-query-advanced/.mvn/wrapper/maven-wrapper.jar diff --git a/kafka-streams/kafka-streams-interactive-query/.mvn/wrapper/maven-wrapper.properties b/kafka-streams/kafka-streams-interactive-query-advanced/.mvn/wrapper/maven-wrapper.properties similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/.mvn/wrapper/maven-wrapper.properties rename to kafka-streams/kafka-streams-interactive-query-advanced/.mvn/wrapper/maven-wrapper.properties diff --git a/kafka-streams/kafka-streams-interactive-query/README.adoc b/kafka-streams/kafka-streams-interactive-query-advanced/README.adoc similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/README.adoc rename to kafka-streams/kafka-streams-interactive-query-advanced/README.adoc diff --git a/kafka-streams/kafka-streams-interactive-query/docker/docker-compose.yml b/kafka-streams/kafka-streams-interactive-query-advanced/docker/docker-compose.yml similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/docker/docker-compose.yml rename to kafka-streams/kafka-streams-interactive-query-advanced/docker/docker-compose.yml diff --git a/kafka-streams/kafka-streams-interactive-query/docker/start-kafka-shell.sh b/kafka-streams/kafka-streams-interactive-query-advanced/docker/start-kafka-shell.sh similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/docker/start-kafka-shell.sh rename to kafka-streams/kafka-streams-interactive-query-advanced/docker/start-kafka-shell.sh diff --git a/kafka-streams/kafka-streams-interactive-query/mvnw b/kafka-streams/kafka-streams-interactive-query-advanced/mvnw similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/mvnw rename to kafka-streams/kafka-streams-interactive-query-advanced/mvnw diff --git a/kafka-streams/kafka-streams-interactive-query/mvnw.cmd b/kafka-streams/kafka-streams-interactive-query-advanced/mvnw.cmd similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/mvnw.cmd rename to kafka-streams/kafka-streams-interactive-query-advanced/mvnw.cmd diff --git a/kafka-streams/kafka-streams-interactive-query/pom.xml b/kafka-streams/kafka-streams-interactive-query-advanced/pom.xml similarity index 97% rename from kafka-streams/kafka-streams-interactive-query/pom.xml rename to kafka-streams/kafka-streams-interactive-query-advanced/pom.xml index 17077fa..b3fef22 100644 --- a/kafka-streams/kafka-streams-interactive-query/pom.xml +++ b/kafka-streams/kafka-streams-interactive-query-advanced/pom.xml @@ -4,11 +4,11 @@ 4.0.0 kafka.streams.interactive.query - kafka-streams-interactive-query + kafka-streams-interactive-query-advanced 0.0.1-SNAPSHOT jar - kafka-streams-interactive-query + kafka-streams-interactive-query-advanced Demo project for Spring Boot diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQuerySample.java b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQuerySample.java similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQuerySample.java rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQuerySample.java diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/java/kafka/streams/interactive/query/Producers.java b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/java/kafka/streams/interactive/query/Producers.java similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/java/kafka/streams/interactive/query/Producers.java rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/java/kafka/streams/interactive/query/Producers.java diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/java/kafka/streams/interactive/query/SongPlayCountBean.java b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/java/kafka/streams/interactive/query/SongPlayCountBean.java similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/java/kafka/streams/interactive/query/SongPlayCountBean.java rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/java/kafka/streams/interactive/query/SongPlayCountBean.java diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/resources/application.yml b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/application.yml similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/resources/application.yml rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/application.yml diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/resources/avro/kafka/streams/interactive/query/playevent.avsc b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/avro/kafka/streams/interactive/query/playevent.avsc similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/resources/avro/kafka/streams/interactive/query/playevent.avsc rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/avro/kafka/streams/interactive/query/playevent.avsc diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/resources/avro/kafka/streams/interactive/query/song.avsc b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/avro/kafka/streams/interactive/query/song.avsc similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/resources/avro/kafka/streams/interactive/query/song.avsc rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/avro/kafka/streams/interactive/query/song.avsc diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/resources/avro/kafka/streams/interactive/query/songplaycount.avsc b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/avro/kafka/streams/interactive/query/songplaycount.avsc similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/resources/avro/kafka/streams/interactive/query/songplaycount.avsc rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/avro/kafka/streams/interactive/query/songplaycount.avsc diff --git a/kafka-streams/kafka-streams-interactive-query/src/main/resources/logback.xml b/kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/logback.xml similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/main/resources/logback.xml rename to kafka-streams/kafka-streams-interactive-query-advanced/src/main/resources/logback.xml diff --git a/kafka-streams/kafka-streams-interactive-query/src/test/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQueryTests.java b/kafka-streams/kafka-streams-interactive-query-advanced/src/test/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQueryTests.java similarity index 100% rename from kafka-streams/kafka-streams-interactive-query/src/test/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQueryTests.java rename to kafka-streams/kafka-streams-interactive-query-advanced/src/test/java/kafka/streams/interactive/query/KafkaStreamsInteractiveQueryTests.java diff --git a/kstream/kstream-interactive-query/.jdk8 b/kafka-streams/kafka-streams-interactive-query-basic/.jdk8 similarity index 100% rename from kstream/kstream-interactive-query/.jdk8 rename to kafka-streams/kafka-streams-interactive-query-basic/.jdk8 diff --git a/kstream/kstream-interactive-query/.mvn/wrapper/maven-wrapper.jar b/kafka-streams/kafka-streams-interactive-query-basic/.mvn/wrapper/maven-wrapper.jar similarity index 100% rename from kstream/kstream-interactive-query/.mvn/wrapper/maven-wrapper.jar rename to kafka-streams/kafka-streams-interactive-query-basic/.mvn/wrapper/maven-wrapper.jar diff --git a/kstream/kstream-interactive-query/.mvn/wrapper/maven-wrapper.properties b/kafka-streams/kafka-streams-interactive-query-basic/.mvn/wrapper/maven-wrapper.properties similarity index 100% rename from kstream/kstream-interactive-query/.mvn/wrapper/maven-wrapper.properties rename to kafka-streams/kafka-streams-interactive-query-basic/.mvn/wrapper/maven-wrapper.properties diff --git a/kstream/kstream-interactive-query/README.adoc b/kafka-streams/kafka-streams-interactive-query-basic/README.adoc similarity index 70% rename from kstream/kstream-interactive-query/README.adoc rename to kafka-streams/kafka-streams-interactive-query-basic/README.adoc index 2a519e3..99259c0 100644 --- a/kstream/kstream-interactive-query/README.adoc +++ b/kafka-streams/kafka-streams-interactive-query-basic/README.adoc @@ -15,7 +15,9 @@ This sample uses lambda expressions and thus requires Java 8+. 3. cd $KAFKA_HOME 4. Start the console producer: + Assuming that you are running kafka on a docker container on mac osx. Change the zookeeper IP address accordingly otherwise. + -`bin/kafka-console-producer.sh --broker-list 192.168.99.100:9092 --topic products` +`bin/kafka-console-producer.sh --broker-list localhost:9092 --topic products` + +`bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --key-deserializer org.apache.kafka.common.serialization.IntegerDeserializer --property print.key=true --topic product-counts print.value=true --value-deserializer org.apache.kafka.common.serialization.LongDeserializer` === Running the app: @@ -23,15 +25,15 @@ Go to the root of the repository and do: `./mvnw clean package` -`java -jar target/kstream-interactive-query-0.0.1-SNAPSHOT.jar --kstream.product.tracker.productIds=123,124,125` +`java -jar target/kafka-streams-interactive-query-basic-0.0.1-SNAPSHOT.jar --app.product.tracker.productIds=123,124,125` The above command will track products with ID's 123,124 and 125 and print their counts seen so far every 30 seconds. * By default we use the docker container IP (mac osx specific) in the `application.yml` for Kafka broker and zookeeper. Change it in `application.yml` (which requires a rebuild) or pass them as runtime arguments as below. -`spring.cloud.stream.kstream.binder.brokers=` + -`spring.cloud.stream.kstream.binder.zkNodes=` +`spring.cloud.stream.kafka.streams.binder.brokers=` + +`spring.cloud.stream.kafka.streams.binder.zkNodes=` Enter the following in the console producer (one line at a time) and watch the output on the console (or IDE) where the application is running. diff --git a/kstream/kstream-interactive-query/docker/docker-compose.yml b/kafka-streams/kafka-streams-interactive-query-basic/docker/docker-compose.yml similarity index 100% rename from kstream/kstream-interactive-query/docker/docker-compose.yml rename to kafka-streams/kafka-streams-interactive-query-basic/docker/docker-compose.yml diff --git a/kstream/kstream-interactive-query/docker/start-kafka-shell.sh b/kafka-streams/kafka-streams-interactive-query-basic/docker/start-kafka-shell.sh similarity index 100% rename from kstream/kstream-interactive-query/docker/start-kafka-shell.sh rename to kafka-streams/kafka-streams-interactive-query-basic/docker/start-kafka-shell.sh diff --git a/kstream/kstream-interactive-query/mvnw b/kafka-streams/kafka-streams-interactive-query-basic/mvnw similarity index 100% rename from kstream/kstream-interactive-query/mvnw rename to kafka-streams/kafka-streams-interactive-query-basic/mvnw diff --git a/kafka-streams/kafka-streams-interactive-query-basic/pom.xml b/kafka-streams/kafka-streams-interactive-query-basic/pom.xml new file mode 100644 index 0000000..042a81a --- /dev/null +++ b/kafka-streams/kafka-streams-interactive-query-basic/pom.xml @@ -0,0 +1,89 @@ + + + 4.0.0 + + kafka.streams.interactive.query + kafka-streams-interactive-query-basic + 0.0.1-SNAPSHOT + jar + + kafka-streams-interactive-query-basic + Spring Cloud Stream sample for KStream interactive queries + + + org.springframework.boot + spring-boot-starter-parent + 2.0.0.BUILD-SNAPSHOT + + + + + UTF-8 + UTF-8 + 1.8 + Finchley.BUILD-SNAPSHOT + + + + + org.springframework.boot + spring-boot-starter + + + org.springframework.kafka + spring-kafka + 2.1.3.RELEASE + + + org.springframework.cloud + spring-cloud-stream-binder-kafka-streams + 2.0.0.BUILD-SNAPSHOT + + + + org.springframework.boot + spring-boot-starter-test + test + + + + org.springframework.boot + spring-boot-starter-web + + + + + + + + org.springframework.cloud + spring-cloud-dependencies + ${spring-cloud.version} + pom + import + + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + + + + spring-milestones + Spring Milestones + http://repo.spring.io/libs-milestone-local + + false + + + + + diff --git a/kstream/kstream-interactive-query/src/main/java/kstream/word/count/kstreamwordcount/KStreamInteractiveQueryApplication.java b/kafka-streams/kafka-streams-interactive-query-basic/src/main/java/kafka/streams/product/tracker/KafkaStreamsInteractiveQueryApplication.java similarity index 83% rename from kstream/kstream-interactive-query/src/main/java/kstream/word/count/kstreamwordcount/KStreamInteractiveQueryApplication.java rename to kafka-streams/kafka-streams-interactive-query-basic/src/main/java/kafka/streams/product/tracker/KafkaStreamsInteractiveQueryApplication.java index ee92bad..da610e6 100644 --- a/kstream/kstream-interactive-query/src/main/java/kstream/word/count/kstreamwordcount/KStreamInteractiveQueryApplication.java +++ b/kafka-streams/kafka-streams-interactive-query-basic/src/main/java/kafka/streams/product/tracker/KafkaStreamsInteractiveQueryApplication.java @@ -14,18 +14,13 @@ * limitations under the License. */ -package kstream.word.count.kstreamwordcount; - -import java.util.Set; -import java.util.stream.Collectors; +package kafka.streams.product.tracker; import org.apache.kafka.common.serialization.Serdes; -import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.state.QueryableStoreTypes; import org.apache.kafka.streams.state.ReadOnlyKeyValueStore; - import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; @@ -34,22 +29,25 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.kstream.annotations.KStreamProcessor; -import org.springframework.kafka.core.KStreamBuilderFactoryBean; +import org.springframework.cloud.stream.binder.kafka.streams.QueryableStoreRegistry; +import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.kafka.support.serializer.JsonSerde; import org.springframework.messaging.handler.annotation.SendTo; import org.springframework.scheduling.annotation.EnableScheduling; import org.springframework.scheduling.annotation.Scheduled; import org.springframework.util.StringUtils; +import java.util.Set; +import java.util.stream.Collectors; + @SpringBootApplication -public class KStreamInteractiveQueryApplication { +public class KafkaStreamsInteractiveQueryApplication { public static void main(String[] args) { - SpringApplication.run(KStreamInteractiveQueryApplication.class, args); + SpringApplication.run(KafkaStreamsInteractiveQueryApplication.class, args); } - @EnableBinding(KStreamProcessor.class) + @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration @EnableConfigurationProperties(ProductTrackerProperties.class) @EnableScheduling @@ -58,7 +56,7 @@ public class KStreamInteractiveQueryApplication { private static final String STORE_NAME = "prod-id-count-store"; @Autowired - private KStreamBuilderFactoryBean kStreamBuilderFactoryBean; + private QueryableStoreRegistry queryableStoreRegistry; @Autowired ProductTrackerProperties productTrackerProperties; @@ -86,8 +84,7 @@ public class KStreamInteractiveQueryApplication { @Scheduled(fixedRate = 30000, initialDelay = 5000) public void printProductCounts() { if (keyValueStore == null) { - KafkaStreams streams = kStreamBuilderFactoryBean.getKafkaStreams(); - keyValueStore = streams.store(STORE_NAME, QueryableStoreTypes.keyValueStore()); + keyValueStore = queryableStoreRegistry.getQueryableStoreType(STORE_NAME, QueryableStoreTypes.keyValueStore()); } for (Integer id : productIds()) { @@ -97,7 +94,7 @@ public class KStreamInteractiveQueryApplication { } - @ConfigurationProperties(prefix = "kstream.product.tracker") + @ConfigurationProperties(prefix = "app.product.tracker") static class ProductTrackerProperties { private String productIds; diff --git a/kstream/kstream-interactive-query/src/main/resources/application.yml b/kafka-streams/kafka-streams-interactive-query-basic/src/main/resources/application.yml similarity index 64% rename from kstream/kstream-interactive-query/src/main/resources/application.yml rename to kafka-streams/kafka-streams-interactive-query-basic/src/main/resources/application.yml index 0fe6361..5ba82d0 100644 --- a/kstream/kstream-interactive-query/src/main/resources/application.yml +++ b/kafka-streams/kafka-streams-interactive-query-basic/src/main/resources/application.yml @@ -1,12 +1,12 @@ spring.cloud.stream.bindings.output.contentType: application/json -spring.cloud.stream.kstream.binder.configuration.commit.interval.ms: 1000 -spring.cloud.stream.kstream: +spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms: 1000 +spring.cloud.stream.kafka.streams: binder.configuration: key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde bindings.output.producer: keySerde: org.apache.kafka.common.serialization.Serdes$IntegerSerde - valueSerde: org.apache.kafka.common.serialization.Serdes$StringSerde + valueSerde: org.apache.kafka.common.serialization.Serdes$LongSerde spring.cloud.stream.bindings.output: destination: product-counts producer: @@ -16,6 +16,6 @@ spring.cloud.stream.bindings.input: destination: products consumer: headerMode: raw -spring.cloud.stream.kstream.binder: - brokers: 192.168.99.100 - zkNodes: 192.168.99.100 \ No newline at end of file +spring.cloud.stream.kafka.streams.binder: + brokers: localhost #192.168.99.100 + zkNodes: localhost #192.168.99.100 \ No newline at end of file diff --git a/kstream/kstream-interactive-query/src/main/resources/logback.xml b/kafka-streams/kafka-streams-interactive-query-basic/src/main/resources/logback.xml similarity index 100% rename from kstream/kstream-interactive-query/src/main/resources/logback.xml rename to kafka-streams/kafka-streams-interactive-query-basic/src/main/resources/logback.xml diff --git a/kstream/kstream-product-tracker/.jdk8 b/kafka-streams/kafka-streams-product-tracker/.jdk8 similarity index 100% rename from kstream/kstream-product-tracker/.jdk8 rename to kafka-streams/kafka-streams-product-tracker/.jdk8 diff --git a/kstream/kstream-product-tracker/.mvn/wrapper/maven-wrapper.jar b/kafka-streams/kafka-streams-product-tracker/.mvn/wrapper/maven-wrapper.jar similarity index 100% rename from kstream/kstream-product-tracker/.mvn/wrapper/maven-wrapper.jar rename to kafka-streams/kafka-streams-product-tracker/.mvn/wrapper/maven-wrapper.jar diff --git a/kstream/kstream-product-tracker/.mvn/wrapper/maven-wrapper.properties b/kafka-streams/kafka-streams-product-tracker/.mvn/wrapper/maven-wrapper.properties similarity index 100% rename from kstream/kstream-product-tracker/.mvn/wrapper/maven-wrapper.properties rename to kafka-streams/kafka-streams-product-tracker/.mvn/wrapper/maven-wrapper.properties diff --git a/kstream/kstream-product-tracker/README.adoc b/kafka-streams/kafka-streams-product-tracker/README.adoc similarity index 84% rename from kstream/kstream-product-tracker/README.adoc rename to kafka-streams/kafka-streams-product-tracker/README.adoc index f2cf4cd..1583348 100644 --- a/kstream/kstream-product-tracker/README.adoc +++ b/kafka-streams/kafka-streams-product-tracker/README.adoc @@ -27,7 +27,7 @@ Go to the root of the repository and do: `./mvnw clean package` -`java -jar target/kstream-product-tracker-0.0.1-SNAPSHOT.jar --kstream.product.tracker.productIds=123,124,125 --spring.cloud.stream.kstream.timeWindow.length=60000 --spring.cloud.stream.kstream.timeWindow.advanceBy=30000 --spring.cloud.stream.bindings.input.destination=products` +`java -jar target/kafka-streams-product-tracker-0.0.1-SNAPSHOT.jar --app.product.tracker.productIds=123,124,125 --spring.cloud.stream.kafka.streams.timeWindow.length=60000 --spring.cloud.stream.kafka.streams.timeWindow.advanceBy=30000` The above command will track products with ID's 123,124 and 125 every 30 seconds with the counts from the last minute. In other words, every 30 seconds a new 1 minute window is started. @@ -52,7 +52,7 @@ Enter the following in the console producer (one line at a time) and watch the o The default time window is configured for 30 seconds and you can change that using the following property. -`kstream.word.count.windowLength` (value is expressed in milliseconds) +`spring.cloud.stream.kafka.streams.timeWindow.length` (value is expressed in milliseconds) -In order to switch to a hopping window, you can use the `kstream.word.count.advanceBy` (value in milliseconds). +In order to switch to a hopping window, you can use the `spring.cloud.stream.kafka.streams.timeWindow.advanceBy` (value in milliseconds). This will create an overlapped hopping windows depending on the value you provide. diff --git a/kstream/kstream-product-tracker/docker/docker-compose.yml b/kafka-streams/kafka-streams-product-tracker/docker/docker-compose.yml similarity index 100% rename from kstream/kstream-product-tracker/docker/docker-compose.yml rename to kafka-streams/kafka-streams-product-tracker/docker/docker-compose.yml diff --git a/kstream/kstream-product-tracker/docker/start-kafka-shell.sh b/kafka-streams/kafka-streams-product-tracker/docker/start-kafka-shell.sh similarity index 100% rename from kstream/kstream-product-tracker/docker/start-kafka-shell.sh rename to kafka-streams/kafka-streams-product-tracker/docker/start-kafka-shell.sh diff --git a/kstream/kstream-product-tracker/mvnw b/kafka-streams/kafka-streams-product-tracker/mvnw similarity index 100% rename from kstream/kstream-product-tracker/mvnw rename to kafka-streams/kafka-streams-product-tracker/mvnw diff --git a/kafka-streams/kafka-streams-product-tracker/pom.xml b/kafka-streams/kafka-streams-product-tracker/pom.xml new file mode 100644 index 0000000..c5315a5 --- /dev/null +++ b/kafka-streams/kafka-streams-product-tracker/pom.xml @@ -0,0 +1,89 @@ + + + 4.0.0 + + kafka.streams.product.tracker + kafka-streams-product-tracker + 0.0.1-SNAPSHOT + jar + + kafka-streams-product-tracker + Demo project for Spring Boot + + + org.springframework.boot + spring-boot-starter-parent + 2.0.0.BUILD-SNAPSHOT + + + + + UTF-8 + UTF-8 + 1.8 + Finchley.BUILD-SNAPSHOT + + + + + org.springframework.boot + spring-boot-starter + + + org.springframework.kafka + spring-kafka + 2.1.3.RELEASE + + + org.springframework.cloud + spring-cloud-stream-binder-kafka-streams + 2.0.0.BUILD-SNAPSHOT + + + + org.springframework.boot + spring-boot-starter-test + test + + + + org.springframework.boot + spring-boot-starter-web + + + + + + + + org.springframework.cloud + spring-cloud-dependencies + ${spring-cloud.version} + pom + import + + + + + + + + org.springframework.boot + spring-boot-maven-plugin + + + + + + + spring-milestones + Spring Milestones + http://repo.spring.io/libs-milestone-local + + false + + + + + diff --git a/kstream/kstream-product-tracker/src/main/java/kstream/word/count/kstreamwordcount/KStreamProductTrackerApplication.java b/kafka-streams/kafka-streams-product-tracker/src/main/java/kafka/streams/product/tracker/KafkaStreamsProductTrackerApplication.java similarity index 91% rename from kstream/kstream-product-tracker/src/main/java/kstream/word/count/kstreamwordcount/KStreamProductTrackerApplication.java rename to kafka-streams/kafka-streams-product-tracker/src/main/java/kafka/streams/product/tracker/KafkaStreamsProductTrackerApplication.java index 9784a03..0304afc 100644 --- a/kstream/kstream-product-tracker/src/main/java/kstream/word/count/kstreamwordcount/KStreamProductTrackerApplication.java +++ b/kafka-streams/kafka-streams-product-tracker/src/main/java/kafka/streams/product/tracker/KafkaStreamsProductTrackerApplication.java @@ -14,18 +14,11 @@ * limitations under the License. */ -package kstream.word.count.kstreamwordcount; - -import java.time.Instant; -import java.time.LocalTime; -import java.time.ZoneId; -import java.util.Set; -import java.util.stream.Collectors; +package kafka.streams.product.tracker; import org.apache.kafka.streams.KeyValue; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.TimeWindows; - import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; @@ -34,19 +27,25 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.kstream.annotations.KStreamProcessor; +import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor; import org.springframework.kafka.support.serializer.JsonSerde; import org.springframework.messaging.handler.annotation.SendTo; import org.springframework.util.StringUtils; +import java.time.Instant; +import java.time.LocalTime; +import java.time.ZoneId; +import java.util.Set; +import java.util.stream.Collectors; + @SpringBootApplication -public class KStreamProductTrackerApplication { +public class KafkaStreamsProductTrackerApplication { public static void main(String[] args) { - SpringApplication.run(KStreamProductTrackerApplication.class, args); + SpringApplication.run(KafkaStreamsProductTrackerApplication.class, args); } - @EnableBinding(KStreamProcessor.class) + @EnableBinding(KafkaStreamsProcessor.class) @EnableAutoConfiguration @EnableConfigurationProperties(ProductTrackerProperties.class) public static class ProductCountApplication { @@ -78,7 +77,7 @@ public class KStreamProductTrackerApplication { } - @ConfigurationProperties(prefix = "kstream.product.tracker") + @ConfigurationProperties(prefix = "app.product.tracker") static class ProductTrackerProperties { private String productIds; diff --git a/kstream/kstream-product-tracker/src/main/resources/application.yml b/kafka-streams/kafka-streams-product-tracker/src/main/resources/application.yml similarity index 67% rename from kstream/kstream-product-tracker/src/main/resources/application.yml rename to kafka-streams/kafka-streams-product-tracker/src/main/resources/application.yml index 5e37c8f..11b9f51 100644 --- a/kstream/kstream-product-tracker/src/main/resources/application.yml +++ b/kafka-streams/kafka-streams-product-tracker/src/main/resources/application.yml @@ -1,6 +1,6 @@ spring.cloud.stream.bindings.output.contentType: application/json -spring.cloud.stream.kstream.binder.configuration.commit.interval.ms: 1000 -spring.cloud.stream.kstream: +spring.cloud.stream.kafka.streams.binder.configuration.commit.interval.ms: 1000 +spring.cloud.stream.kafka.streams: binder.configuration: key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde @@ -10,11 +10,11 @@ spring.cloud.stream.bindings.output: destination: product-counts producer: headerMode: raw - useNativeEncoding: true + #useNativeEncoding: true spring.cloud.stream.bindings.input: destination: products consumer: headerMode: raw -spring.cloud.stream.kstream.binder: - brokers: 192.168.99.100 - zkNodes: 192.168.99.100 \ No newline at end of file +spring.cloud.stream.kafka.streams.binder: + brokers: localhost #192.168.99.100 + zkNodes: localhost #192.168.99.100 \ No newline at end of file diff --git a/kstream/kstream-product-tracker/src/main/resources/logback.xml b/kafka-streams/kafka-streams-product-tracker/src/main/resources/logback.xml similarity index 100% rename from kstream/kstream-product-tracker/src/main/resources/logback.xml rename to kafka-streams/kafka-streams-product-tracker/src/main/resources/logback.xml diff --git a/kafka-streams/pom.xml b/kafka-streams/pom.xml index 9c42dcf..eeb7395 100644 --- a/kafka-streams/pom.xml +++ b/kafka-streams/pom.xml @@ -16,10 +16,13 @@ kafka-streams-word-count - kafka-streams-branching-sample + kafka-streams-branching kafka-streams-dlq-sample kafka-streams-table-join - kafka-streams-interactive-query - kafka-streams-message-channel + kafka-streams-interactive-query-basic + kafka-streams-interactive-query-advanced + kafka-streams-message-channel + kafka-streams-product-tracker + kafka-streams-aggregate diff --git a/kstream/kstream-interactive-query/pom.xml b/kstream/kstream-interactive-query/pom.xml deleted file mode 100644 index 36c1346..0000000 --- a/kstream/kstream-interactive-query/pom.xml +++ /dev/null @@ -1,53 +0,0 @@ - - - 4.0.0 - - kstream.product.tracker - kstream-interactive-query - 0.0.1-SNAPSHOT - jar - - kstream-interactive-query - Spring Cloud Stream sample for KStream interactive queries - - - org.springframework.cloud - spring-cloud-stream-sample-kstream-parent - 1.2.0.BUILD-SNAPSHOT - - - - UTF-8 - UTF-8 - 1.8 - - - - - org.springframework.cloud - spring-cloud-stream-binder-kstream - - - - - - - org.springframework.boot - spring-boot-maven-plugin - - - - - - - spring-milestones - Spring Milestones - http://repo.spring.io/libs-milestone-local - - false - - - - - diff --git a/kstream/kstream-product-tracker/pom.xml b/kstream/kstream-product-tracker/pom.xml deleted file mode 100644 index 37f3241..0000000 --- a/kstream/kstream-product-tracker/pom.xml +++ /dev/null @@ -1,53 +0,0 @@ - - - 4.0.0 - - kstream.product.tracker - kstream-product-tracker - 0.0.1-SNAPSHOT - jar - - kstream-product-tracker - Demo project for Spring Boot - - - org.springframework.cloud - spring-cloud-stream-sample-kstream-parent - 1.2.0.BUILD-SNAPSHOT - - - - UTF-8 - UTF-8 - 1.8 - - - - - org.springframework.cloud - spring-cloud-stream-binder-kstream - - - - - - - org.springframework.boot - spring-boot-maven-plugin - - - - - - - spring-milestones - Spring Milestones - http://repo.spring.io/libs-milestone-local - - false - - - - - diff --git a/kstream/kstream-word-count/.gitignore b/kstream/kstream-word-count/.gitignore deleted file mode 100644 index 2af7cef..0000000 --- a/kstream/kstream-word-count/.gitignore +++ /dev/null @@ -1,24 +0,0 @@ -target/ -!.mvn/wrapper/maven-wrapper.jar - -### STS ### -.apt_generated -.classpath -.factorypath -.project -.settings -.springBeans - -### IntelliJ IDEA ### -.idea -*.iws -*.iml -*.ipr - -### NetBeans ### -nbproject/private/ -build/ -nbbuild/ -dist/ -nbdist/ -.nb-gradle/ \ No newline at end of file diff --git a/kstream/kstream-word-count/.jdk8 b/kstream/kstream-word-count/.jdk8 deleted file mode 100644 index e69de29..0000000 diff --git a/kstream/kstream-word-count/.mvn/wrapper/maven-wrapper.jar b/kstream/kstream-word-count/.mvn/wrapper/maven-wrapper.jar deleted file mode 100644 index 9cc84ea..0000000 Binary files a/kstream/kstream-word-count/.mvn/wrapper/maven-wrapper.jar and /dev/null differ diff --git a/kstream/kstream-word-count/.mvn/wrapper/maven-wrapper.properties b/kstream/kstream-word-count/.mvn/wrapper/maven-wrapper.properties deleted file mode 100644 index c315043..0000000 --- a/kstream/kstream-word-count/.mvn/wrapper/maven-wrapper.properties +++ /dev/null @@ -1 +0,0 @@ -distributionUrl=https://repo1.maven.org/maven2/org/apache/maven/apache-maven/3.5.0/apache-maven-3.5.0-bin.zip diff --git a/kstream/kstream-word-count/README.adoc b/kstream/kstream-word-count/README.adoc deleted file mode 100644 index c4b38d7..0000000 --- a/kstream/kstream-word-count/README.adoc +++ /dev/null @@ -1,43 +0,0 @@ -== What is this app? - -This is an example of a Spring Cloud Stream processor using Kafka Streams support. - -The example is based on the word count application from the https://github.com/confluentinc/examples/blob/3.2.x/kafka-streams/src/main/java/io/confluent/examples/streams/WordCountLambdaExample.java[reference documentation]. -It uses a single input and a single output. -In essence, the application receives text messages from an input topic and computes word occurrence counts in a configurable time window and report that in an output topic. -This sample uses lambda expressions and thus requires Java 8+. - -==== Starting Kafka in a docker container - -* Skip steps 1-3 if you already have a non-Docker Kafka environment. - -1. Go to the docker directory in this repo and invoke the command `docker-compose up -d`. -2. Ensure that in the docker directory and then invoke the script `start-kafka-shell.sh` -3. cd $KAFKA_HOME -4. Start the console producer: + -Assuming that you are running kafka on a docker container on mac osx. Change the zookeeper IP address accordingly otherwise. + -`bin/kafka-console-producer.sh --broker-list 192.168.99.100:9092 --topic words` -5. Start the console consumer: + -Assuming that you are running kafka on a docker container on mac osx. Change the zookeeper IP address accordingly otherwise. + -`bin/kafka-console-consumer.sh --bootstrap-server 192.168.99.100:9092 --topic counts` - -=== Running the app: - -Go to the root of the repository and do: `./mvnw clean package` - -`java -jar target/kstream-word-count-0.0.1-SNAPSHOT.jar --spring.cloud.stream.kstream.timeWindow.length=5000` - -* By default we use the docker container IP (mac osx specific) in the `application.yml` for Kafka broker and zookeeper. -Change it in `application.yml` (which requires a rebuild) or pass them as runtime arguments as below. - -`spring.cloud.stream.kstream.binder.brokers=` + -`spring.cloud.stream.kstream.binder.zkNodes=` - -Enter some text in the console producer and watch the output in the console consumer. - -In order to switch to a hopping window, you can use the `spring.cloud.stream.kstream.timeWindow.advanceBy` (value in milliseconds). -This will create an overlapped hopping windows depending on the value you provide. - -Here is an example with 2 overlapping windows (window length of 10 seconds and a hop (advance) by 5 seconds: - -`java -jar target/kstream-word-count-0.0.1-SNAPSHOT.jar --spring.cloud.stream.kstream.timeWindow.length=10000 --spring.cloud.stream.kstream.timeWindow.advanceBy=5000` \ No newline at end of file diff --git a/kstream/kstream-word-count/docker/docker-compose.yml b/kstream/kstream-word-count/docker/docker-compose.yml deleted file mode 100644 index ecf41fb..0000000 --- a/kstream/kstream-word-count/docker/docker-compose.yml +++ /dev/null @@ -1,14 +0,0 @@ -version: '2' -services: - zookeeper: - image: wurstmeister/zookeeper - ports: - - "2181:2181" - kafka: - image: wurstmeister/kafka - ports: - - "9092:9092" - environment: - KAFKA_ADVERTISED_HOST_NAME: 192.168.99.100 - KAFKA_ADVERTISED_PORT: 9092 - KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181 diff --git a/kstream/kstream-word-count/docker/start-kafka-shell.sh b/kstream/kstream-word-count/docker/start-kafka-shell.sh deleted file mode 100755 index 62663e4..0000000 --- a/kstream/kstream-word-count/docker/start-kafka-shell.sh +++ /dev/null @@ -1,2 +0,0 @@ -#!/bin/bash -docker run --rm -v /var/run/docker.sock:/var/run/docker.sock -e HOST_IP=$1 -e ZK=$2 -i -t wurstmeister/kafka /bin/bash diff --git a/kstream/kstream-word-count/mvnw b/kstream/kstream-word-count/mvnw deleted file mode 100755 index 5bf251c..0000000 --- a/kstream/kstream-word-count/mvnw +++ /dev/null @@ -1,225 +0,0 @@ -#!/bin/sh -# ---------------------------------------------------------------------------- -# Licensed to the Apache Software Foundation (ASF) under one -# or more contributor license agreements. See the NOTICE file -# distributed with this work for additional information -# regarding copyright ownership. The ASF licenses this file -# to you 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. -# ---------------------------------------------------------------------------- - -# ---------------------------------------------------------------------------- -# Maven2 Start Up Batch script -# -# Required ENV vars: -# ------------------ -# JAVA_HOME - location of a JDK home dir -# -# Optional ENV vars -# ----------------- -# M2_HOME - location of maven2's installed home dir -# MAVEN_OPTS - parameters passed to the Java VM when running Maven -# e.g. to debug Maven itself, use -# set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000 -# MAVEN_SKIP_RC - flag to disable loading of mavenrc files -# ---------------------------------------------------------------------------- - -if [ -z "$MAVEN_SKIP_RC" ] ; then - - if [ -f /etc/mavenrc ] ; then - . /etc/mavenrc - fi - - if [ -f "$HOME/.mavenrc" ] ; then - . "$HOME/.mavenrc" - fi - -fi - -# OS specific support. $var _must_ be set to either true or false. -cygwin=false; -darwin=false; -mingw=false -case "`uname`" in - CYGWIN*) cygwin=true ;; - MINGW*) mingw=true;; - Darwin*) darwin=true - # Use /usr/libexec/java_home if available, otherwise fall back to /Library/Java/Home - # See https://developer.apple.com/library/mac/qa/qa1170/_index.html - if [ -z "$JAVA_HOME" ]; then - if [ -x "/usr/libexec/java_home" ]; then - export JAVA_HOME="`/usr/libexec/java_home`" - else - export JAVA_HOME="/Library/Java/Home" - fi - fi - ;; -esac - -if [ -z "$JAVA_HOME" ] ; then - if [ -r /etc/gentoo-release ] ; then - JAVA_HOME=`java-config --jre-home` - fi -fi - -if [ -z "$M2_HOME" ] ; then - ## resolve links - $0 may be a link to maven's home - PRG="$0" - - # need this for relative symlinks - while [ -h "$PRG" ] ; do - ls=`ls -ld "$PRG"` - link=`expr "$ls" : '.*-> \(.*\)$'` - if expr "$link" : '/.*' > /dev/null; then - PRG="$link" - else - PRG="`dirname "$PRG"`/$link" - fi - done - - saveddir=`pwd` - - M2_HOME=`dirname "$PRG"`/.. - - # make it fully qualified - M2_HOME=`cd "$M2_HOME" && pwd` - - cd "$saveddir" - # echo Using m2 at $M2_HOME -fi - -# For Cygwin, ensure paths are in UNIX format before anything is touched -if $cygwin ; then - [ -n "$M2_HOME" ] && - M2_HOME=`cygpath --unix "$M2_HOME"` - [ -n "$JAVA_HOME" ] && - JAVA_HOME=`cygpath --unix "$JAVA_HOME"` - [ -n "$CLASSPATH" ] && - CLASSPATH=`cygpath --path --unix "$CLASSPATH"` -fi - -# For Migwn, ensure paths are in UNIX format before anything is touched -if $mingw ; then - [ -n "$M2_HOME" ] && - M2_HOME="`(cd "$M2_HOME"; pwd)`" - [ -n "$JAVA_HOME" ] && - JAVA_HOME="`(cd "$JAVA_HOME"; pwd)`" - # TODO classpath? -fi - -if [ -z "$JAVA_HOME" ]; then - javaExecutable="`which javac`" - if [ -n "$javaExecutable" ] && ! [ "`expr \"$javaExecutable\" : '\([^ ]*\)'`" = "no" ]; then - # readlink(1) is not available as standard on Solaris 10. - readLink=`which readlink` - if [ ! `expr "$readLink" : '\([^ ]*\)'` = "no" ]; then - if $darwin ; then - javaHome="`dirname \"$javaExecutable\"`" - javaExecutable="`cd \"$javaHome\" && pwd -P`/javac" - else - javaExecutable="`readlink -f \"$javaExecutable\"`" - fi - javaHome="`dirname \"$javaExecutable\"`" - javaHome=`expr "$javaHome" : '\(.*\)/bin'` - JAVA_HOME="$javaHome" - export JAVA_HOME - fi - fi -fi - -if [ -z "$JAVACMD" ] ; then - if [ -n "$JAVA_HOME" ] ; then - if [ -x "$JAVA_HOME/jre/sh/java" ] ; then - # IBM's JDK on AIX uses strange locations for the executables - JAVACMD="$JAVA_HOME/jre/sh/java" - else - JAVACMD="$JAVA_HOME/bin/java" - fi - else - JAVACMD="`which java`" - fi -fi - -if [ ! -x "$JAVACMD" ] ; then - echo "Error: JAVA_HOME is not defined correctly." >&2 - echo " We cannot execute $JAVACMD" >&2 - exit 1 -fi - -if [ -z "$JAVA_HOME" ] ; then - echo "Warning: JAVA_HOME environment variable is not set." -fi - -CLASSWORLDS_LAUNCHER=org.codehaus.plexus.classworlds.launcher.Launcher - -# traverses directory structure from process work directory to filesystem root -# first directory with .mvn subdirectory is considered project base directory -find_maven_basedir() { - - if [ -z "$1" ] - then - echo "Path not specified to find_maven_basedir" - return 1 - fi - - basedir="$1" - wdir="$1" - while [ "$wdir" != '/' ] ; do - if [ -d "$wdir"/.mvn ] ; then - basedir=$wdir - break - fi - # workaround for JBEAP-8937 (on Solaris 10/Sparc) - if [ -d "${wdir}" ]; then - wdir=`cd "$wdir/.."; pwd` - fi - # end of workaround - done - echo "${basedir}" -} - -# concatenates all lines of a file -concat_lines() { - if [ -f "$1" ]; then - echo "$(tr -s '\n' ' ' < "$1")" - fi -} - -BASE_DIR=`find_maven_basedir "$(pwd)"` -if [ -z "$BASE_DIR" ]; then - exit 1; -fi - -export MAVEN_PROJECTBASEDIR=${MAVEN_BASEDIR:-"$BASE_DIR"} -echo $MAVEN_PROJECTBASEDIR -MAVEN_OPTS="$(concat_lines "$MAVEN_PROJECTBASEDIR/.mvn/jvm.config") $MAVEN_OPTS" - -# For Cygwin, switch paths to Windows format before running java -if $cygwin; then - [ -n "$M2_HOME" ] && - M2_HOME=`cygpath --path --windows "$M2_HOME"` - [ -n "$JAVA_HOME" ] && - JAVA_HOME=`cygpath --path --windows "$JAVA_HOME"` - [ -n "$CLASSPATH" ] && - CLASSPATH=`cygpath --path --windows "$CLASSPATH"` - [ -n "$MAVEN_PROJECTBASEDIR" ] && - MAVEN_PROJECTBASEDIR=`cygpath --path --windows "$MAVEN_PROJECTBASEDIR"` -fi - -WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain - -exec "$JAVACMD" \ - $MAVEN_OPTS \ - -classpath "$MAVEN_PROJECTBASEDIR/.mvn/wrapper/maven-wrapper.jar" \ - "-Dmaven.home=${M2_HOME}" "-Dmaven.multiModuleProjectDirectory=${MAVEN_PROJECTBASEDIR}" \ - ${WRAPPER_LAUNCHER} $MAVEN_CONFIG "$@" diff --git a/kstream/kstream-word-count/mvnw.cmd b/kstream/kstream-word-count/mvnw.cmd deleted file mode 100644 index 019bd74..0000000 --- a/kstream/kstream-word-count/mvnw.cmd +++ /dev/null @@ -1,143 +0,0 @@ -@REM ---------------------------------------------------------------------------- -@REM Licensed to the Apache Software Foundation (ASF) under one -@REM or more contributor license agreements. See the NOTICE file -@REM distributed with this work for additional information -@REM regarding copyright ownership. The ASF licenses this file -@REM to you under the Apache License, Version 2.0 (the -@REM "License"); you may not use this file except in compliance -@REM with the License. You may obtain a copy of the License at -@REM -@REM http://www.apache.org/licenses/LICENSE-2.0 -@REM -@REM Unless required by applicable law or agreed to in writing, -@REM software distributed under the License is distributed on an -@REM "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY -@REM KIND, either express or implied. See the License for the -@REM specific language governing permissions and limitations -@REM under the License. -@REM ---------------------------------------------------------------------------- - -@REM ---------------------------------------------------------------------------- -@REM Maven2 Start Up Batch script -@REM -@REM Required ENV vars: -@REM JAVA_HOME - location of a JDK home dir -@REM -@REM Optional ENV vars -@REM M2_HOME - location of maven2's installed home dir -@REM MAVEN_BATCH_ECHO - set to 'on' to enable the echoing of the batch commands -@REM MAVEN_BATCH_PAUSE - set to 'on' to wait for a key stroke before ending -@REM MAVEN_OPTS - parameters passed to the Java VM when running Maven -@REM e.g. to debug Maven itself, use -@REM set MAVEN_OPTS=-Xdebug -Xrunjdwp:transport=dt_socket,server=y,suspend=y,address=8000 -@REM MAVEN_SKIP_RC - flag to disable loading of mavenrc files -@REM ---------------------------------------------------------------------------- - -@REM Begin all REM lines with '@' in case MAVEN_BATCH_ECHO is 'on' -@echo off -@REM enable echoing my setting MAVEN_BATCH_ECHO to 'on' -@if "%MAVEN_BATCH_ECHO%" == "on" echo %MAVEN_BATCH_ECHO% - -@REM set %HOME% to equivalent of $HOME -if "%HOME%" == "" (set "HOME=%HOMEDRIVE%%HOMEPATH%") - -@REM Execute a user defined script before this one -if not "%MAVEN_SKIP_RC%" == "" goto skipRcPre -@REM check for pre script, once with legacy .bat ending and once with .cmd ending -if exist "%HOME%\mavenrc_pre.bat" call "%HOME%\mavenrc_pre.bat" -if exist "%HOME%\mavenrc_pre.cmd" call "%HOME%\mavenrc_pre.cmd" -:skipRcPre - -@setlocal - -set ERROR_CODE=0 - -@REM To isolate internal variables from possible post scripts, we use another setlocal -@setlocal - -@REM ==== START VALIDATION ==== -if not "%JAVA_HOME%" == "" goto OkJHome - -echo. -echo Error: JAVA_HOME not found in your environment. >&2 -echo Please set the JAVA_HOME variable in your environment to match the >&2 -echo location of your Java installation. >&2 -echo. -goto error - -:OkJHome -if exist "%JAVA_HOME%\bin\java.exe" goto init - -echo. -echo Error: JAVA_HOME is set to an invalid directory. >&2 -echo JAVA_HOME = "%JAVA_HOME%" >&2 -echo Please set the JAVA_HOME variable in your environment to match the >&2 -echo location of your Java installation. >&2 -echo. -goto error - -@REM ==== END VALIDATION ==== - -:init - -@REM Find the project base dir, i.e. the directory that contains the folder ".mvn". -@REM Fallback to current working directory if not found. - -set MAVEN_PROJECTBASEDIR=%MAVEN_BASEDIR% -IF NOT "%MAVEN_PROJECTBASEDIR%"=="" goto endDetectBaseDir - -set EXEC_DIR=%CD% -set WDIR=%EXEC_DIR% -:findBaseDir -IF EXIST "%WDIR%"\.mvn goto baseDirFound -cd .. -IF "%WDIR%"=="%CD%" goto baseDirNotFound -set WDIR=%CD% -goto findBaseDir - -:baseDirFound -set MAVEN_PROJECTBASEDIR=%WDIR% -cd "%EXEC_DIR%" -goto endDetectBaseDir - -:baseDirNotFound -set MAVEN_PROJECTBASEDIR=%EXEC_DIR% -cd "%EXEC_DIR%" - -:endDetectBaseDir - -IF NOT EXIST "%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config" goto endReadAdditionalConfig - -@setlocal EnableExtensions EnableDelayedExpansion -for /F "usebackq delims=" %%a in ("%MAVEN_PROJECTBASEDIR%\.mvn\jvm.config") do set JVM_CONFIG_MAVEN_PROPS=!JVM_CONFIG_MAVEN_PROPS! %%a -@endlocal & set JVM_CONFIG_MAVEN_PROPS=%JVM_CONFIG_MAVEN_PROPS% - -:endReadAdditionalConfig - -SET MAVEN_JAVA_EXE="%JAVA_HOME%\bin\java.exe" - -set WRAPPER_JAR="%MAVEN_PROJECTBASEDIR%\.mvn\wrapper\maven-wrapper.jar" -set WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain - -%MAVEN_JAVA_EXE% %JVM_CONFIG_MAVEN_PROPS% %MAVEN_OPTS% %MAVEN_DEBUG_OPTS% -classpath %WRAPPER_JAR% "-Dmaven.multiModuleProjectDirectory=%MAVEN_PROJECTBASEDIR%" %WRAPPER_LAUNCHER% %MAVEN_CONFIG% %* -if ERRORLEVEL 1 goto error -goto end - -:error -set ERROR_CODE=1 - -:end -@endlocal & set ERROR_CODE=%ERROR_CODE% - -if not "%MAVEN_SKIP_RC%" == "" goto skipRcPost -@REM check for post script, once with legacy .bat ending and once with .cmd ending -if exist "%HOME%\mavenrc_post.bat" call "%HOME%\mavenrc_post.bat" -if exist "%HOME%\mavenrc_post.cmd" call "%HOME%\mavenrc_post.cmd" -:skipRcPost - -@REM pause the script if MAVEN_BATCH_PAUSE is set to 'on' -if "%MAVEN_BATCH_PAUSE%" == "on" pause - -if "%MAVEN_TERMINATE_CMD%" == "on" exit %ERROR_CODE% - -exit /B %ERROR_CODE% diff --git a/kstream/kstream-word-count/pom.xml b/kstream/kstream-word-count/pom.xml deleted file mode 100644 index 319bef6..0000000 --- a/kstream/kstream-word-count/pom.xml +++ /dev/null @@ -1,53 +0,0 @@ - - - 4.0.0 - - kstream.word.count - kstream-word-count - 0.0.1-SNAPSHOT - jar - - kstream-word-count - Demo project for Spring Boot - - - org.springframework.cloud - spring-cloud-stream-sample-kstream-parent - 1.2.0.BUILD-SNAPSHOT - - - - UTF-8 - UTF-8 - 1.8 - - - - - org.springframework.cloud - spring-cloud-stream-binder-kstream - - - - - - - org.springframework.boot - spring-boot-maven-plugin - - - - - - - spring-milestones - Spring Milestones - http://repo.spring.io/libs-milestone-local - - false - - - - - diff --git a/kstream/kstream-word-count/src/main/java/kstream/word/count/kstreamwordcount/KStreamWordCountApplication.java b/kstream/kstream-word-count/src/main/java/kstream/word/count/kstreamwordcount/KStreamWordCountApplication.java deleted file mode 100644 index 8ed098a..0000000 --- a/kstream/kstream-word-count/src/main/java/kstream/word/count/kstreamwordcount/KStreamWordCountApplication.java +++ /dev/null @@ -1,114 +0,0 @@ -/* - * Copyright 2017 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 kstream.word.count.kstreamwordcount; - -import java.util.Arrays; -import java.util.Date; - -import org.apache.kafka.common.serialization.Serdes; -import org.apache.kafka.streams.KeyValue; -import org.apache.kafka.streams.kstream.KStream; -import org.apache.kafka.streams.kstream.TimeWindows; - -import org.springframework.beans.factory.annotation.Autowired; -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.EnableAutoConfiguration; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.binder.kstream.annotations.KStreamProcessor; -import org.springframework.messaging.handler.annotation.SendTo; - -@SpringBootApplication -public class KStreamWordCountApplication { - - public static void main(String[] args) { - SpringApplication.run(KStreamWordCountApplication.class, args); - } - - @EnableBinding(KStreamProcessor.class) - @EnableAutoConfiguration - public static class WordCountProcessorApplication { - - @Autowired - private TimeWindows timeWindows; - - @StreamListener("input") - @SendTo("output") - public KStream process(KStream input) { - - return input - .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+"))) - .map((key, value) -> new KeyValue<>(value, value)) - .groupByKey(Serdes.String(), Serdes.String()) - .count(timeWindows, "WordCounts") - .toStream() - .map((key, value) -> new KeyValue<>(null, new WordCount(key.key(), value, new Date(key.window().start()), new Date(key.window().end())))); - } - - } - - static class WordCount { - - private String word; - - private long count; - - private Date start; - - private Date end; - - WordCount(String word, long count, Date start, Date end) { - this.word = word; - this.count = count; - this.start = start; - this.end = end; - } - - public String getWord() { - return word; - } - - public void setWord(String word) { - this.word = word; - } - - public long getCount() { - return count; - } - - public void setCount(long count) { - this.count = count; - } - - public Date getStart() { - return start; - } - - public void setStart(Date start) { - this.start = start; - } - - public Date getEnd() { - return end; - } - - public void setEnd(Date end) { - this.end = end; - } - } -} diff --git a/kstream/kstream-word-count/src/main/resources/application.yml b/kstream/kstream-word-count/src/main/resources/application.yml deleted file mode 100644 index 0688d67..0000000 --- a/kstream/kstream-word-count/src/main/resources/application.yml +++ /dev/null @@ -1,17 +0,0 @@ -spring.cloud.stream.bindings.output.contentType: application/json -spring.cloud.stream.kstream.binder.configuration.commit.interval.ms: 1000 -spring.cloud.stream.kstream.binder.configuration: - key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde - value.serde: org.apache.kafka.common.serialization.Serdes$StringSerde -spring.cloud.stream.bindings.output: - destination: counts - producer: - headerMode: raw - useNativeEncoding: true -spring.cloud.stream.bindings.input: - destination: words - consumer: - headerMode: raw -spring.cloud.stream.kstream.binder: - brokers: 192.168.99.100 - zkNodes: 192.168.99.100 \ No newline at end of file diff --git a/kstream/kstream-word-count/src/main/resources/logback.xml b/kstream/kstream-word-count/src/main/resources/logback.xml deleted file mode 100644 index 870ac9e..0000000 --- a/kstream/kstream-word-count/src/main/resources/logback.xml +++ /dev/null @@ -1,12 +0,0 @@ - - - - - %d{ISO8601} %5p %t %c{2}:%L - %m%n - - - - - - - \ No newline at end of file diff --git a/kstream/pom.xml b/kstream/pom.xml deleted file mode 100644 index ccbc734..0000000 --- a/kstream/pom.xml +++ /dev/null @@ -1,22 +0,0 @@ - - - 4.0.0 - - spring-cloud-stream-sample-kstream-parent - pom - - spring-cloud-stream-sample-kstream-parent - Parent Project for KStream Samples - - - org.springframework.cloud - spring-cloud-stream-samples - 1.2.0.BUILD-SNAPSHOT - - - - kstream-word-count - kstream-product-tracker - kstream-interactive-query - - diff --git a/mvnw b/mvnw index e351148..5bf251c 100755 --- a/mvnw +++ b/mvnw @@ -54,38 +54,16 @@ case "`uname`" in CYGWIN*) cygwin=true ;; MINGW*) mingw=true;; Darwin*) darwin=true - # - # Look for the Apple JDKs first to preserve the existing behaviour, and then look - # for the new JDKs provided by Oracle. - # - if [ -z "$JAVA_HOME" ] && [ -L /System/Library/Frameworks/JavaVM.framework/Versions/CurrentJDK ] ; then - # - # Apple JDKs - # - export JAVA_HOME=/System/Library/Frameworks/JavaVM.framework/Versions/CurrentJDK/Home - fi - - if [ -z "$JAVA_HOME" ] && [ -L /System/Library/Java/JavaVirtualMachines/CurrentJDK ] ; then - # - # Apple JDKs - # - export JAVA_HOME=/System/Library/Java/JavaVirtualMachines/CurrentJDK/Contents/Home - fi - - if [ -z "$JAVA_HOME" ] && [ -L "/Library/Java/JavaVirtualMachines/CurrentJDK" ] ; then - # - # Oracle JDKs - # - export JAVA_HOME=/Library/Java/JavaVirtualMachines/CurrentJDK/Contents/Home - fi - - if [ -z "$JAVA_HOME" ] && [ -x "/usr/libexec/java_home" ]; then - # - # Apple JDKs - # - export JAVA_HOME=`/usr/libexec/java_home` - fi - ;; + # Use /usr/libexec/java_home if available, otherwise fall back to /Library/Java/Home + # See https://developer.apple.com/library/mac/qa/qa1170/_index.html + if [ -z "$JAVA_HOME" ]; then + if [ -x "/usr/libexec/java_home" ]; then + export JAVA_HOME="`/usr/libexec/java_home`" + else + export JAVA_HOME="/Library/Java/Home" + fi + fi + ;; esac if [ -z "$JAVA_HOME" ] ; then @@ -184,27 +162,28 @@ fi CLASSWORLDS_LAUNCHER=org.codehaus.plexus.classworlds.launcher.Launcher -# For Cygwin, switch paths to Windows format before running java -if $cygwin; then - [ -n "$M2_HOME" ] && - M2_HOME=`cygpath --path --windows "$M2_HOME"` - [ -n "$JAVA_HOME" ] && - JAVA_HOME=`cygpath --path --windows "$JAVA_HOME"` - [ -n "$CLASSPATH" ] && - CLASSPATH=`cygpath --path --windows "$CLASSPATH"` -fi - # traverses directory structure from process work directory to filesystem root # first directory with .mvn subdirectory is considered project base directory find_maven_basedir() { - local basedir=$(pwd) - local wdir=$(pwd) + + if [ -z "$1" ] + then + echo "Path not specified to find_maven_basedir" + return 1 + fi + + basedir="$1" + wdir="$1" while [ "$wdir" != '/' ] ; do if [ -d "$wdir"/.mvn ] ; then basedir=$wdir break fi - wdir=$(cd "$wdir/.."; pwd) + # workaround for JBEAP-8937 (on Solaris 10/Sparc) + if [ -d "${wdir}" ]; then + wdir=`cd "$wdir/.."; pwd` + fi + # end of workaround done echo "${basedir}" } @@ -216,13 +195,26 @@ concat_lines() { fi } -export MAVEN_PROJECTBASEDIR=${MAVEN_BASEDIR:-$(find_maven_basedir)} +BASE_DIR=`find_maven_basedir "$(pwd)"` +if [ -z "$BASE_DIR" ]; then + exit 1; +fi + +export MAVEN_PROJECTBASEDIR=${MAVEN_BASEDIR:-"$BASE_DIR"} +echo $MAVEN_PROJECTBASEDIR MAVEN_OPTS="$(concat_lines "$MAVEN_PROJECTBASEDIR/.mvn/jvm.config") $MAVEN_OPTS" -# Provide a "standardized" way to retrieve the CLI args that will -# work with both Windows and non-Windows executions. -MAVEN_CMD_LINE_ARGS="$MAVEN_CONFIG $@" -export MAVEN_CMD_LINE_ARGS +# For Cygwin, switch paths to Windows format before running java +if $cygwin; then + [ -n "$M2_HOME" ] && + M2_HOME=`cygpath --path --windows "$M2_HOME"` + [ -n "$JAVA_HOME" ] && + JAVA_HOME=`cygpath --path --windows "$JAVA_HOME"` + [ -n "$CLASSPATH" ] && + CLASSPATH=`cygpath --path --windows "$CLASSPATH"` + [ -n "$MAVEN_PROJECTBASEDIR" ] && + MAVEN_PROJECTBASEDIR=`cygpath --path --windows "$MAVEN_PROJECTBASEDIR"` +fi WRAPPER_LAUNCHER=org.apache.maven.wrapper.MavenWrapperMain @@ -230,5 +222,4 @@ exec "$JAVACMD" \ $MAVEN_OPTS \ -classpath "$MAVEN_PROJECTBASEDIR/.mvn/wrapper/maven-wrapper.jar" \ "-Dmaven.home=${M2_HOME}" "-Dmaven.multiModuleProjectDirectory=${MAVEN_PROJECTBASEDIR}" \ - ${WRAPPER_LAUNCHER} "$@" - + ${WRAPPER_LAUNCHER} $MAVEN_CONFIG "$@" diff --git a/pom.xml b/pom.xml index e77d353..39e73d2 100644 --- a/pom.xml +++ b/pom.xml @@ -29,14 +29,13 @@ non-self-contained-aggregate-app multibinder multibinder-differentsystems - rxjava-processor multi-io stream-listener reactive-processor-kafka test-embedded-kafka kinesis-produce-consume - kstream testing + kafka-streams diff --git a/rxjava-processor/pom.xml b/rxjava-processor/pom.xml deleted file mode 100644 index 432c239..0000000 --- a/rxjava-processor/pom.xml +++ /dev/null @@ -1,71 +0,0 @@ - - - 4.0.0 - - spring-cloud-stream-sample-rxjava - jar - - spring-cloud-stream-sample-rxjava - Demo project for RxJava module - - - org.springframework.cloud - spring-cloud-stream-samples - 1.2.0.BUILD-SNAPSHOT - - - - UTF-8 - demo.RxJavaApplication - 1.8 - - - - - org.springframework.cloud - spring-cloud-stream-rxjava - - - org.springframework.cloud - spring-cloud-stream-binder-rabbit - - - org.springframework.cloud - spring-cloud-stream-reactive - - - io.projectreactor - reactor-core - 3.0.5.RELEASE - - - io.reactivex - rxjava - 1.1.10 - - - org.springframework.boot - spring-boot-configuration-processor - true - - - - org.springframework.boot - spring-boot-starter-test - test - - - - - - - org.springframework.boot - spring-boot-maven-plugin - - exec - - - - - - diff --git a/rxjava-processor/src/main/java/demo/RxJavaApplication.java b/rxjava-processor/src/main/java/demo/RxJavaApplication.java deleted file mode 100644 index 78fc214..0000000 --- a/rxjava-processor/src/main/java/demo/RxJavaApplication.java +++ /dev/null @@ -1,33 +0,0 @@ -/* - * Copyright 2015 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 demo; - -import org.springframework.boot.SpringApplication; -import org.springframework.boot.autoconfigure.SpringBootApplication; -import org.springframework.context.annotation.ComponentScan; - -/** - * @author Ilayaperumal Gopinathan - */ -@SpringBootApplication -public class RxJavaApplication { - - public static void main(String[] args) { - SpringApplication.run(RxJavaApplication.class, args); - } - -} diff --git a/rxjava-processor/src/main/java/demo/RxJavaTransformer.java b/rxjava-processor/src/main/java/demo/RxJavaTransformer.java deleted file mode 100644 index e39633e..0000000 --- a/rxjava-processor/src/main/java/demo/RxJavaTransformer.java +++ /dev/null @@ -1,54 +0,0 @@ -/* - * Copyright 2015 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 demo; - -import java.util.List; - -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import rx.Observable; - -import org.springframework.cloud.stream.annotation.EnableBinding; -import org.springframework.cloud.stream.annotation.StreamListener; -import org.springframework.cloud.stream.messaging.Processor; -import org.springframework.messaging.handler.annotation.SendTo; - -@EnableBinding(Processor.class) -public class RxJavaTransformer { - - private static Logger logger = LoggerFactory.getLogger(RxJavaTransformer.class); - - @StreamListener(Processor.INPUT) - @SendTo(Processor.OUTPUT) - public Observable processor(Observable inputStream) { - return inputStream.map(data -> { - logger.info("Got data = " + data); - return data; - }).buffer(5).map(data -> String.valueOf(avg(data))); - } - - private static Double avg(List data) { - double sum = 0; - double count = 0; - for(String d : data) { - count++; - sum += Double.valueOf(d); - } - return sum/count; - } - -} diff --git a/rxjava-processor/src/main/resources/application.yml b/rxjava-processor/src/main/resources/application.yml deleted file mode 100644 index 9b51624..0000000 --- a/rxjava-processor/src/main/resources/application.yml +++ /dev/null @@ -1,10 +0,0 @@ -server: - port: 8082 -spring: - cloud: - stream: - bindings: - output: - destination: xformed - input: - destination: testtock diff --git a/rxjava-processor/src/test/java/demo/ModuleApplicationTests.java b/rxjava-processor/src/test/java/demo/ModuleApplicationTests.java deleted file mode 100644 index 937640b..0000000 --- a/rxjava-processor/src/test/java/demo/ModuleApplicationTests.java +++ /dev/null @@ -1,38 +0,0 @@ -/* - * Copyright 2015 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 demo; - -import org.junit.Ignore; -import org.junit.Test; -import org.junit.runner.RunWith; -import org.springframework.boot.test.context.SpringBootTest; -import org.springframework.test.annotation.DirtiesContext; -import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; -import org.springframework.test.context.web.WebAppConfiguration; - -@RunWith(SpringJUnit4ClassRunner.class) -@SpringBootTest(classes = RxJavaApplication.class) -@WebAppConfiguration -@DirtiesContext -public class ModuleApplicationTests { - - @Test - @Ignore - public void contextLoads() { - } - -}