cleanup/polishing
This commit is contained in:
@@ -22,35 +22,35 @@ import org.springframework.boot.SpringApplication;
|
||||
import org.springframework.boot.autoconfigure.SpringBootApplication;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Import;
|
||||
import org.springframework.integration.support.MutableMessage;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.messaging.MessageHeaders;
|
||||
import org.springframework.messaging.support.MessageBuilder;
|
||||
|
||||
@SpringBootApplication
|
||||
@Import({ org.springframework.cloud.fn.supplier.http.HttpSupplierConfiguration.class })
|
||||
public class HttpIngest {
|
||||
public class HttpClickIngest {
|
||||
|
||||
@Bean
|
||||
public Function<Message<?>, Message<?>> byteArrayToLong() {
|
||||
public Function<Message<?>, Message<UserClicks>> toUserClicks() {
|
||||
return message -> {
|
||||
if (message.getPayload() instanceof byte[]) {
|
||||
MessageHeaders headers = message.getHeaders();
|
||||
String contentType = headers.containsKey("contentType") ?
|
||||
headers.get("contentType").toString() :
|
||||
"application/json";
|
||||
if (contentType.contains("text") || contentType.contains("json") || contentType
|
||||
.contains("x-spring-tuple")) {
|
||||
message = MessageBuilder
|
||||
.withPayload(Long.valueOf(new String((byte[]) ((byte[]) message.getPayload()))))
|
||||
.copyHeaders(message.getHeaders()).build();
|
||||
}
|
||||
}
|
||||
|
||||
return message;
|
||||
return new MutableMessage<>(new UserClicks(
|
||||
(String) message.getHeaders().get("kafka_receivedMessageKey"),
|
||||
Long.valueOf(new String((byte[])message.getPayload()))));
|
||||
};
|
||||
}
|
||||
|
||||
public class UserClicks {
|
||||
|
||||
private String username;
|
||||
|
||||
private Long clicks;
|
||||
|
||||
public UserClicks(String username, Long clicks) {
|
||||
this.username = username;
|
||||
this.clicks = clicks;
|
||||
}
|
||||
}
|
||||
|
||||
public static void main(String[] args) {
|
||||
SpringApplication.run(HttpIngest.class, args);
|
||||
SpringApplication.run(HttpClickIngest.class, args);
|
||||
}
|
||||
}
|
||||
@@ -28,7 +28,7 @@ import org.springframework.messaging.support.MessageBuilder;
|
||||
|
||||
@SpringBootApplication
|
||||
@Import({ org.springframework.cloud.fn.supplier.http.HttpSupplierConfiguration.class })
|
||||
public class HttpIngest {
|
||||
public class HttpRegionIngest {
|
||||
|
||||
@Bean
|
||||
public Function<Message<?>, Message<?>> byteArrayToString() {
|
||||
@@ -41,7 +41,7 @@ public class HttpIngest {
|
||||
if (contentType.contains("text") || contentType.contains("json") || contentType
|
||||
.contains("x-spring-tuple")) {
|
||||
message = MessageBuilder
|
||||
.withPayload(new String((byte[]) ((byte[]) message.getPayload())))
|
||||
.withPayload(new String((byte[]) (message.getPayload())))
|
||||
.copyHeaders(message.getHeaders()).build();
|
||||
}
|
||||
}
|
||||
@@ -51,6 +51,6 @@ public class HttpIngest {
|
||||
}
|
||||
|
||||
public static void main(String[] args) {
|
||||
SpringApplication.run(HttpIngest.class, args);
|
||||
SpringApplication.run(HttpRegionIngest.class, args);
|
||||
}
|
||||
}
|
||||
@@ -1,4 +1,3 @@
|
||||
configuration-properties.classes=org.springframework.cloud.fn.supplier.http.HttpSupplierProperties,\
|
||||
org.springframework.cloud.fn.supplier.http.HttpSupplierProperties$Cors
|
||||
configuration-properties.names=server.port
|
||||
configuration-properties.outbound-ports=clicks
|
||||
|
||||
Reference in New Issue
Block a user