From b1025b0c6c5b9a909e5ca56f5497058f31c52884 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 18 Sep 2019 16:30:25 -0400 Subject: [PATCH] Schema registry samples update - interation 2 Updating kafka-streams-schema-evolution with Confluent Scheam registry samples --- .../custom-confluent-producer1/pom.xml | 82 +++++++++++++++-- .../{FooSerde.java => FooSerializer.java} | 2 +- .../producer1/Producer1Application.java | 30 ++++--- .../src/main/resources/application.yml | 2 +- .../{FooSerde.java => FooSerializer.java} | 2 +- .../producer2/Producer2Application.java | 30 ++++--- .../src/main/resources/application.yml | 2 +- .../kafka-streams-confluent-consumer/pom.xml | 89 ++++++++++++++++--- .../consumer/CountVersionApplication.java | 14 +-- .../src/main/resources/application.yml | 6 +- 10 files changed, 205 insertions(+), 54 deletions(-) rename schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer1/src/main/java/sample/producer1/{FooSerde.java => FooSerializer.java} (89%) rename schema-registry-samples/kafka-streams-schema-evolution/custom-confluent-producer2/src/main/java/sample/producer2/{FooSerde.java => FooSerializer.java} (89%) 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 input) { + @Bean + public Function, KStream> process() { + + //The following Serde definitions are not needed in the topoloyy below + //as we are not using it. However, if your topoloyg explicitly uses this + //Serde, you need to configure this with the schema registry url as below. final Map serdeConfig = Collections.singletonMap( AbstractKafkaAvroSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://localhost:8081"); @@ -55,7 +59,7 @@ public class CountVersionApplication { final SpecificAvroSerde sensorSerde = new SpecificAvroSerde<>(); sensorSerde.configure(serdeConfig, false); - return input + return input -> input .map((k, value) -> { String newKey = "v1"; if (value.getId().toString().endsWith("v2")) { diff --git a/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/resources/application.yml b/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/resources/application.yml index bb21ff1..082f3da 100644 --- a/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/resources/application.yml +++ b/schema-registry-samples/kafka-streams-schema-evolution/kafka-streams-confluent-consumer/src/main/resources/application.yml @@ -1,11 +1,9 @@ server.port: 9998 spring.application.name: kafka-streams-confluent-consumer -spring.cloud.stream.bindings.output: +spring.cloud.stream.bindings.process_out: destination: sensor-versions -spring.cloud.stream.bindings.input: +spring.cloud.stream.bindings.process_in: destination: sensors - consumer: - useNativeDecoding: true spring.cloud.stream.kafka.streams.binder: brokers: localhost configuration: