Divide one record to multiple in deserialize stage

Viewed 179

I've tried to get Kinesis data by Flink.

In my case, there are multiple message in one record. How can I divide it to multiple records? (I'll send it to Elasticsearch.)

I tried to search it but I couldn't find appropriate answer.

What my code does is get data from Kinesis, decompress gzip, convert it to string, then use objectMapper.readvalue to my POJO for java.

There's two POJO's: one for whole event, one for LogEvents.

{
  "messageType":"DATA_MESSAGE","owner":"<account id>",
  "logGroup":"<clustername>","logStream":"<log stream name>",
  "subscriptionFilters":["<subscription name>"],
  "logEvents":[
    {"id":"<id>","timestamp":<timestamp>,"message":"msg 1"},
    {"id":"<id>","timestamp":<timestamp>,"message":"msg 2"},
    {"id":"<id>","timestamp":<timestamp>,"message":"msg 3"},
    {"id":"<id>","timestamp":<timestamp>,"message":"msg 4"},
  ]
}
3 Answers

Your processElement() method which de-serializes the data can call output() multiple times. Each time with a different element of the batch. I.e., you need to loop over the elements in the batch calling output() inside the loop.

The next operator in the chain will get individual elements.

  • Maybe you can use the flatmap function

  • The core method of the FlatMapFunction. Takes an element from the input data set and transforms it into zero, one, or more elements.

      datastream.flatMap(
              new FlatMapFunction<POJO, LogEvent>() {
                  @Override
                  public void flatMap(POJO input, Collector<LogEvent> out) throws Exception {
                      LogEvent logEvent = xxxxx;
                      out.collect(logEvent);
                  }
              });
    

The example JSON has array attributes. You need to write a custom user-defined function to convert JSON into multiple rows and use the flatMap function in Flink to break nesting.

By doing this you will be able to extract rows from the JSON in the format you need and send them as rows to the subsequent Flink operators.

Related