From be81f87c3cb93108d24e69ac2f76d2660ad637af Mon Sep 17 00:00:00 2001 From: DomWos Date: Mon, 3 Sep 2018 12:08:45 +0200 Subject: [PATCH] Switched fields to be static to remove initialization in method. --- .../main/java/com/baeldung/flink/FlinkDataPipeline.java | 2 +- .../baeldung/flink/schema/BackupSerializationSchema.java | 6 +++++- .../flink/schema/InputMessageDeserializationSchema.java | 9 +++++---- 3 files changed, 11 insertions(+), 6 deletions(-) diff --git a/libraries/src/main/java/com/baeldung/flink/FlinkDataPipeline.java b/libraries/src/main/java/com/baeldung/flink/FlinkDataPipeline.java index 423637bf53..d02b1bcb83 100644 --- a/libraries/src/main/java/com/baeldung/flink/FlinkDataPipeline.java +++ b/libraries/src/main/java/com/baeldung/flink/FlinkDataPipeline.java @@ -48,7 +48,7 @@ public static void createBackup () throws Exception { String inputTopic = "flink_input"; String outputTopic = "flink_output"; String consumerGroup = "baeldung"; - String kafkaAddress = "192.168.99.100:9092"; + String kafkaAddress = "localhost:9092"; StreamExecutionEnvironment environment = StreamExecutionEnvironment.getExecutionEnvironment(); diff --git a/libraries/src/main/java/com/baeldung/flink/schema/BackupSerializationSchema.java b/libraries/src/main/java/com/baeldung/flink/schema/BackupSerializationSchema.java index 4db9556d8d..967b266bb6 100644 --- a/libraries/src/main/java/com/baeldung/flink/schema/BackupSerializationSchema.java +++ b/libraries/src/main/java/com/baeldung/flink/schema/BackupSerializationSchema.java @@ -1,6 +1,8 @@ package com.baeldung.flink.schema; import com.baeldung.flink.model.Backup; +import com.fasterxml.jackson.annotation.JsonAutoDetect; +import com.fasterxml.jackson.annotation.PropertyAccessor; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import org.apache.flink.api.common.serialization.SerializationSchema; @@ -10,12 +12,14 @@ import org.slf4j.LoggerFactory; public class BackupSerializationSchema implements SerializationSchema { - ObjectMapper objectMapper; + static ObjectMapper objectMapper = new ObjectMapper().registerModule(new JavaTimeModule()); + Logger logger = LoggerFactory.getLogger(BackupSerializationSchema.class); @Override public byte[] serialize(Backup backupMessage) { if(objectMapper == null) { + objectMapper.setVisibility(PropertyAccessor.FIELD, JsonAutoDetect.Visibility.ANY); objectMapper = new ObjectMapper().registerModule(new JavaTimeModule()); } try { diff --git a/libraries/src/main/java/com/baeldung/flink/schema/InputMessageDeserializationSchema.java b/libraries/src/main/java/com/baeldung/flink/schema/InputMessageDeserializationSchema.java index 3c81b67ec1..1df456bbe5 100644 --- a/libraries/src/main/java/com/baeldung/flink/schema/InputMessageDeserializationSchema.java +++ b/libraries/src/main/java/com/baeldung/flink/schema/InputMessageDeserializationSchema.java @@ -1,6 +1,8 @@ package com.baeldung.flink.schema; import com.baeldung.flink.model.InputMessage; +import com.fasterxml.jackson.annotation.JsonAutoDetect; +import com.fasterxml.jackson.annotation.PropertyAccessor; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule; import org.apache.flink.api.common.serialization.DeserializationSchema; @@ -11,13 +13,12 @@ import java.io.IOException; public class InputMessageDeserializationSchema implements DeserializationSchema { - ObjectMapper objectMapper; + static ObjectMapper objectMapper = new ObjectMapper().registerModule(new JavaTimeModule()); + @Override public InputMessage deserialize(byte[] bytes) throws IOException { - if(objectMapper == null) { - objectMapper = new ObjectMapper().registerModule(new JavaTimeModule()); - } + return objectMapper.readValue(bytes, InputMessage.class); }