Migrated to the new approach for stream - waiting for runtime message polling mechanism

This commit is contained in:
Marcin Grzejszczak
2020-01-30 10:21:26 +01:00
parent c5f8d880c5
commit d1bc8d8d37
6 changed files with 50 additions and 33 deletions

View File

@@ -18,11 +18,7 @@ package com.example.fraud;
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.Configuration;
@Configuration
@SpringBootApplication
public class Application {

View File

@@ -16,20 +16,27 @@
package com.example.fraud;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.publisher.EmitterProcessor;
import org.springframework.cloud.stream.messaging.Source;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.stereotype.Component;
@Component
class MessagePoller {
private final Source source;
MessagePoller(Source source) {
this.source = source;
private static final Logger log = LoggerFactory.getLogger(MessagePoller.class);
private final EmitterProcessor<String> emitterProcessor;
MessagePoller(EmitterProcessor<String> emitterProcessor) {
this.emitterProcessor = emitterProcessor;
}
public void poll() {
this.source.output().send(MessageBuilder
.withPayload("{\"id\":\"99\",\"temperature\":\"123.45\"}").build());
log.info("Emitting the message");
this.emitterProcessor.onNext("{\"id\":\"99\",\"temperature\":\"123.45\"}");
}
}

View File

@@ -16,15 +16,23 @@
package com.example.fraud;
import org.springframework.cloud.stream.annotation.Output;
import org.springframework.cloud.stream.messaging.Sink;
import org.springframework.messaging.MessageChannel;
import java.util.function.Supplier;
interface MyProcessor extends Sink {
import reactor.core.publisher.EmitterProcessor;
String MY_OUTPUT = "my_output";
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
@Output("my_output")
MessageChannel output();
@Configuration
class MessagingConfiguration {
@Bean
EmitterProcessor<String> sensorDataEmitter() {
return EmitterProcessor.create();
}
@Bean(name = "sensor-data")
Supplier<String> sensorData() {
return () -> "{\"id\":\"99\",\"temperature\":\"123.45\"}";
}
}

View File

@@ -22,6 +22,7 @@ import java.net.URISyntaxException;
import java.net.URL;
import java.nio.file.Files;
import java.util.Arrays;
import java.util.function.Function;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -31,19 +32,16 @@ import org.springframework.cloud.stream.messaging.Sink;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.stereotype.Component;
@Component
class MyProcessorListener {
@Component("my_output")
class MyProcessorListener implements Function<byte[], byte[]> {
private static final Logger log = LoggerFactory.getLogger(MyProcessorListener.class);
private final MyProcessor processor;
private final byte[] expectedInput;
private final byte[] expectedOutput;
MyProcessorListener(MyProcessor processor) {
this.processor = processor;
MyProcessorListener() {
this.expectedInput = forFile("/contracts/messaging/input.pdf");
this.expectedOutput = forFile("/contracts/messaging/output.pdf");
}
@@ -58,15 +56,13 @@ class MyProcessorListener {
}
}
@StreamListener(Sink.INPUT)
void listen(byte[] payload) {
@Override
public byte[] apply(byte[] payload) {
log.info("Got the message!");
if (!Arrays.equals(payload, this.expectedInput)) {
log.error("Input payload size is [" + payload.length + "] and the expected one is [" + this.expectedInput.length + "]");
throw new IllegalStateException("Wrong input");
}
this.processor.output()
.send(MessageBuilder.withPayload(this.expectedOutput).build());
return this.expectedOutput;
}
}

View File

@@ -1,6 +1,7 @@
spring.cloud.stream.bindings.output.contentType=application/json
spring.cloud.stream.bindings.input.contentType=application/octet-stream
spring.cloud.stream.bindings.input.destination=bytes_input
spring.cloud.stream.bindings.my_output.contentType=application/octet-stream
spring.cloud.stream.bindings.my_output.destination=bytes_output
spring.cloud.function.definition=my_output;sensor-data
spring.cloud.stream.bindings.sensor-data-out-0.contentType=application/json
spring.cloud.stream.bindings.my_output-in-0.contentType=application/octet-stream
spring.cloud.stream.bindings.my_output-in-0.destination=bytes_input
spring.cloud.stream.bindings.my_output-out-0.contentType=application/octet-stream
spring.cloud.stream.bindings.my_output-out-0.destination=bytes_output
server.port=0

View File

@@ -21,8 +21,11 @@ import org.junit.Before;
import org.junit.runner.RunWith;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.ImportAutoConfiguration;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.cloud.contract.verifier.messaging.boot.AutoConfigureMessageVerifier;
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
import org.springframework.context.annotation.Configuration;
import org.springframework.test.context.junit4.SpringRunner;
import org.springframework.web.context.WebApplicationContext;
@@ -34,7 +37,7 @@ import org.springframework.web.context.WebApplicationContext;
* @author Marius Bogoevici
*/
@RunWith(SpringRunner.class)
@SpringBootTest(classes = Application.class, properties = "spring.cloud.stream.bindings.output.destination=sensor-data")
@SpringBootTest(classes = {MessagingBase.Config.class, Application.class})
@AutoConfigureMessageVerifier
public abstract class MessagingBase {
@@ -53,4 +56,10 @@ public abstract class MessagingBase {
poller.poll();
}
@Configuration
@ImportAutoConfiguration(TestChannelBinderConfiguration.class)
static class Config {
}
}