Refactor test-embedded-kafka
This commit is contained in:
@@ -16,29 +16,26 @@
|
||||
|
||||
package demo;
|
||||
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.springframework.boot.SpringApplication;
|
||||
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.messaging.Processor;
|
||||
import org.springframework.messaging.handler.annotation.SendTo;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
|
||||
/**
|
||||
* @author Gary Russell
|
||||
*
|
||||
*/
|
||||
@SpringBootApplication
|
||||
@EnableBinding(Processor.class)
|
||||
public class EmbeddedKafkaApplication {
|
||||
|
||||
public static void main(String[] args) {
|
||||
SpringApplication.run(EmbeddedKafkaApplication.class, args);
|
||||
}
|
||||
|
||||
@StreamListener(Processor.INPUT)
|
||||
@SendTo(Processor.OUTPUT)
|
||||
public byte[] handle(byte[] in){
|
||||
return new String(in).toUpperCase().getBytes();
|
||||
@Bean
|
||||
public Function<byte[], byte[]> handle(){
|
||||
return in -> new String(in).toUpperCase().getBytes();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# Binding properties
|
||||
spring.cloud.stream.bindings.output.destination=testEmbeddedOut
|
||||
spring.cloud.stream.bindings.input.destination=testEmbeddedIn
|
||||
spring.cloud.stream.bindings.input.group=embeddedKafkaApplication
|
||||
spring.cloud.stream.bindings.handle-out-0.destination=testEmbeddedOut
|
||||
spring.cloud.stream.bindings.handle-in-0.destination=testEmbeddedIn
|
||||
spring.cloud.stream.bindings.handle-in-0.group=embeddedKafkaApplication
|
||||
|
||||
Reference in New Issue
Block a user