Convert reactive processor to functional model
This commit is contained in:
@@ -18,7 +18,7 @@ The following instructions assume that you are running Kafka as a Docker image.
|
||||
|
||||
* `./mvnw clean package`
|
||||
|
||||
* `java -jar target/reactive-processor-0.0.1-SNAPSHOT.jar`
|
||||
* `java -jar target/reactive-processor-0.0.1-SNAPSHOT-kafka.jar`
|
||||
|
||||
The main application contains the reactive processor that receives textual data for a duration and aggregates them.
|
||||
It then sends the aggregated data through the outbound destination of the processor.
|
||||
@@ -47,6 +47,6 @@ All the instructions above apply here also, but instead of running the default `
|
||||
|
||||
* `./mvnw clean package -P rabbit-binder`
|
||||
|
||||
* `java -jar target/reactive-processor-0.0.1-SNAPSHOT.jar`
|
||||
* `java -jar target/reactive-processor-0.0.1-SNAPSHOT-rabbit.jar`
|
||||
|
||||
Once you are done testing: `docker-compose -f docker-compose-rabbit.yml down`
|
||||
|
||||
@@ -20,6 +20,7 @@ import reactor.core.publisher.Flux;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.function.Function;
|
||||
|
||||
@SpringBootApplication
|
||||
@EnableBinding(Processor.class)
|
||||
@@ -29,12 +30,21 @@ public class ReactiveProcessorApplication {
|
||||
SpringApplication.run(ReactiveProcessorApplication.class, args);
|
||||
}
|
||||
|
||||
@StreamListener
|
||||
@Output(Processor.OUTPUT)
|
||||
public Flux<String> aggregate(@Input(Processor.INPUT) Flux<String> inbound) {
|
||||
return inbound.
|
||||
// @StreamListener
|
||||
// @Output(Processor.OUTPUT)
|
||||
// public Flux<String> aggregate(@Input(Processor.INPUT) Flux<String> inbound) {
|
||||
// return inbound.
|
||||
// log()
|
||||
// .window(Duration.ofSeconds(5), Duration.ofSeconds(5))
|
||||
// .flatMap(w -> w.reduce("", (s1,s2)->s1+s2))
|
||||
// .log();
|
||||
// }
|
||||
|
||||
@Bean
|
||||
public Function<Flux<String>, Flux<String>> aggregate() {
|
||||
return inbound -> inbound.
|
||||
log()
|
||||
.window(Duration.ofSeconds(5), Duration.ofSeconds(5))
|
||||
.window(Duration.ofSeconds(30), Duration.ofSeconds(5))
|
||||
.flatMap(w -> w.reduce("", (s1,s2)->s1+s2))
|
||||
.log();
|
||||
}
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
spring:
|
||||
cloud:
|
||||
stream:
|
||||
function.definition: aggregate
|
||||
bindings:
|
||||
output:
|
||||
destination: transformed
|
||||
|
||||
Reference in New Issue
Block a user