diff --git a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/pom.xml b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/pom.xml
index 1ce12f4..d0a9549 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/pom.xml
+++ b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/pom.xml
@@ -9,17 +9,31 @@
Custom Confluent Producer1
- io.spring.cloud.stream.sample
- spring-cloud-stream-samples-parent
- 0.0.1-SNAPSHOT
- ../../..
+ org.springframework.boot
+ spring-boot-starter-parent
+ 2.2.0.BUILD-SNAPSHOT
+
- 4.0.0
1.8.2
+ 4.0.0
+ Hoxton.BUILD-SNAPSHOT
+
+
+
+ org.springframework.cloud
+ spring-cloud-dependencies
+ ${spring-cloud.version}
+ pom
+ import
+
+
+
+
+
org.springframework.cloud
@@ -97,9 +111,65 @@
+
+ spring-snapshots
+ Spring Snapshots
+ https://repo.spring.io/libs-snapshot-local
+
+ true
+
+
+ false
+
+
+
+ spring-milestones
+ Spring Milestones
+ https://repo.spring.io/libs-milestone-local
+
+ false
+
+
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/release
+
+ false
+
+
confluent
https://packages.confluent.io/maven/
-
+
+
+ spring-snapshots
+ Spring Snapshots
+ https://repo.spring.io/libs-snapshot-local
+
+ true
+
+
+ false
+
+
+
+ spring-milestones
+ Spring Milestones
+ https://repo.spring.io/libs-milestone-local
+
+ false
+
+
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/libs-release-local
+
+ false
+
+
+
+
\ No newline at end of file
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/FooSerde.java b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/FooSerializer.java
similarity index 89%
rename from schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/FooSerde.java
rename to schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/FooSerializer.java
index e4e06e3..4e2630d 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/FooSerde.java
+++ b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/FooSerializer.java
@@ -10,7 +10,7 @@ import java.util.Map;
/**
* @author Soby Chacko
*/
-public class FooSerde extends SpecificAvroSerializer {
+public class FooSerializer extends SpecificAvroSerializer {
@Override
public void configure(Map serializerConfig, boolean isSerializerForRecordKeys) {
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/Producer1Application.java b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/Producer1Application.java
index 2f2acfe..290aa3f 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/Producer1Application.java
+++ b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/Producer1Application.java
@@ -6,6 +6,7 @@ import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.messaging.Source;
+import org.springframework.context.annotation.Bean;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestMethod;
@@ -13,27 +14,22 @@ import org.springframework.web.bind.annotation.RestController;
import java.util.Random;
import java.util.UUID;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.function.Supplier;
@SpringBootApplication
-@EnableBinding(Source.class)
@RestController
-public class Producer1Application {
-
- @Autowired
- private Source source;
+public class Producer1Application{
private Random random = new Random();
+ BlockingQueue unbounded = new LinkedBlockingQueue<>();
+
public static void main(String[] args) {
SpringApplication.run(Producer1Application.class, args);
}
- @RequestMapping(value = "/messages", method = RequestMethod.POST)
- public String sendMessage() {
- source.output().send(MessageBuilder.withPayload(randomSensor()).build());
- return "ok, have fun with v1 payload!";
- }
-
private Sensor randomSensor() {
Sensor sensor = new Sensor();
sensor.setId(UUID.randomUUID().toString() + "-v1");
@@ -42,7 +38,19 @@ public class Producer1Application {
sensor.setTemperature(random.nextFloat() * 50);
return sensor;
}
+
+ @RequestMapping(value = "/messages", method = RequestMethod.POST)
+ public String sendMessage() {
+ unbounded.offer(randomSensor());
+ return "ok, have fun with v1 payload!";
+ }
+
+ @Bean
+ public Supplier supplier() {
+ return () -> unbounded.poll();
+ }
}
+
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/resources/application.yml b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/resources/application.yml
index 768fa2c..54266e4 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/resources/application.yml
+++ b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/resources/application.yml
@@ -10,4 +10,4 @@ server.port: 9009
spring.cloud.stream.kafka.binder.configuration:
schema.registry.url: http://localhost:8081
key.serializer: org.apache.kafka.common.serialization.ByteArraySerializer
- value.serializer: sample.producer1.FooSerde
\ No newline at end of file
+ value.serializer: sample.producer1.FooSerializer
\ No newline at end of file
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/FooSerde.java b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/FooSerializer.java
similarity index 89%
rename from schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/FooSerde.java
rename to schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/FooSerializer.java
index 3e52fb3..d955b34 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/FooSerde.java
+++ b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/FooSerializer.java
@@ -10,7 +10,7 @@ import java.util.Map;
/**
* @author Soby Chacko
*/
-public class FooSerde extends SpecificAvroSerializer {
+public class FooSerializer extends SpecificAvroSerializer {
@Override
public void configure(Map serializerConfig, boolean isSerializerForRecordKeys) {
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/Producer2Application.java b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/Producer2Application.java
index cedc8e9..28c9ff8 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/Producer2Application.java
+++ b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/Producer2Application.java
@@ -6,6 +6,7 @@ import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.messaging.Source;
+import org.springframework.context.annotation.Bean;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RequestMethod;
@@ -13,27 +14,22 @@ import org.springframework.web.bind.annotation.RestController;
import java.util.Random;
import java.util.UUID;
+import java.util.concurrent.BlockingQueue;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.function.Supplier;
@SpringBootApplication
-@EnableBinding(Source.class)
@RestController
-public class Producer2Application {
-
- @Autowired
- private Source source;
+public class Producer2Application{
private Random random = new Random();
+ BlockingQueue unbounded = new LinkedBlockingQueue<>();
+
public static void main(String[] args) {
SpringApplication.run(Producer2Application.class, args);
}
- @RequestMapping(value = "/messages", method = RequestMethod.POST)
- public String sendMessage() {
- source.output().send(MessageBuilder.withPayload(randomSensor()).build());
- return "ok, have fun with v2 payload!";
- }
-
private Sensor randomSensor() {
Sensor sensor = new Sensor();
sensor.setId(UUID.randomUUID().toString() + "-v2");
@@ -44,5 +40,17 @@ public class Producer2Application {
sensor.setMagneticField(null);
return sensor;
}
+
+ @RequestMapping(value = "/messages", method = RequestMethod.POST)
+ public String sendMessage() {
+ unbounded.offer(randomSensor());
+ return "ok, have fun with v2 payload!";
+ }
+
+ @Bean
+ public Supplier supplier() {
+ return () -> unbounded.poll();
+ }
}
+
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/resources/application.yml b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/resources/application.yml
index 631502a..1bbd4a2 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/resources/application.yml
+++ b/schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/resources/application.yml
@@ -10,4 +10,4 @@ server.port: 9010
spring.cloud.stream.kafka.binder.configuration:
schema.registry.url: http://localhost:8081
key.serializer: org.apache.kafka.common.serialization.ByteArraySerializer
- value.serializer: sample.producer2.FooSerde
\ No newline at end of file
+ value.serializer: sample.producer2.FooSerializer
\ No newline at end of file
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/pom.xml b/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/pom.xml
index bc82bd7..3542762 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/pom.xml
+++ b/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/pom.xml
@@ -1,6 +1,5 @@
-
+
4.0.0
kafka-streams-confluent-consumer
@@ -10,22 +9,31 @@
Kafka Streams Confluent Consumer
- io.spring.cloud.stream.sample
- spring-cloud-stream-samples-parent
- 0.0.1-SNAPSHOT
- ../../..
+ org.springframework.boot
+ spring-boot-starter-parent
+ 2.2.0.BUILD-SNAPSHOT
+
- 4.0.0
1.8.2
+ Hoxton.BUILD-SNAPSHOT
+ 4.0.0
+
+
+
+ org.springframework.cloud
+ spring-cloud-dependencies
+ ${spring-cloud.version}
+ pom
+ import
+
+
+
+
-
- org.springframework.cloud
- spring-cloud-stream-schema
-
org.apache.avro
avro
@@ -35,7 +43,6 @@
org.springframework.cloud
spring-cloud-stream-binder-kafka-streams
-
io.confluent
kafka-streams-avro-serde
@@ -87,9 +94,65 @@
+
+ spring-snapshots
+ Spring Snapshots
+ https://repo.spring.io/libs-snapshot-local
+
+ true
+
+
+ false
+
+
+
+ spring-milestones
+ Spring Milestones
+ https://repo.spring.io/libs-milestone-local
+
+ false
+
+
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/release
+
+ false
+
+
confluent
https://packages.confluent.io/maven/
-
+
+
+ spring-snapshots
+ Spring Snapshots
+ https://repo.spring.io/libs-snapshot-local
+
+ true
+
+
+ false
+
+
+
+ spring-milestones
+ Spring Milestones
+ https://repo.spring.io/libs-milestone-local
+
+ false
+
+
+
+ spring-releases
+ Spring Releases
+ https://repo.spring.io/libs-release-local
+
+ false
+
+
+
+
\ No newline at end of file
diff --git a/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/java/sample/consumer/CountVersionApplication.java b/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/java/sample/consumer/CountVersionApplication.java
index 75bcf91..e3fc629 100644
--- a/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/java/sample/consumer/CountVersionApplication.java
+++ b/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/java/sample/consumer/CountVersionApplication.java
@@ -20,15 +20,16 @@ import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.binder.kafka.streams.InteractiveQueryService;
import org.springframework.cloud.stream.binder.kafka.streams.annotations.KafkaStreamsProcessor;
+import org.springframework.context.annotation.Bean;
import org.springframework.messaging.handler.annotation.SendTo;
import org.springframework.scheduling.annotation.EnableScheduling;
import org.springframework.scheduling.annotation.Scheduled;
import java.util.Collections;
import java.util.Map;
+import java.util.function.Function;
@SpringBootApplication
-@EnableBinding(KafkaStreamsProcessor.class)
@EnableScheduling
public class CountVersionApplication {
@@ -45,9 +46,12 @@ public class CountVersionApplication {
SpringApplication.run(CountVersionApplication.class, args);
}
- @StreamListener("input")
- @SendTo("output")
- public KStream process(KStream