Is it safe for a Flink application to have multiple data/key streams in s job all sharing the same Kafka source and sink?

Viewed 213

enter image description here (Goal Updated) My goal on each data stream is:

  • filter different msgs
  • have different event time defined window session gaps
  • consumer from topic and produce to another topic

A fan-out -> fan-in like DAG.

var fanoutStreamOne = new StreamComponents(/*filter, flatmap, etc*/);
var fanoutStreamTwo = new StreamComponents(/*filter, flatmap, etc*/);
var fanoutStreamThree = new StreamComponents(/*filter, flatmap, etc*/);
var fanoutStreams = Set.of(fanoutStreamOne, fanoutStreamTwo, fanoutStreamThree)
var source = new FlinkKafkaConsumer<>(...);
var sink = new FlinkKafkaProducer<>(...);

// creates streams from same source to same sink (Using union())
new streamingJob(source, sink, fanoutStreams).execute();

I am just curious if this affects recovery/checkpoints or performance of the Flink application.

Has anyone had success with this implementation?

And should I have the watermark strategy up front before filtering?

Thank in advance!

1 Answers

Okay, the differenced time gaps are not possible, I think so. I tried it a year ago, with flink 1.7 , and I can't do it. The watermark is global to the application.

To the other problems, if you are using Kafka, yo can read from some topics using regex, and get the topic using the properly deserialization schema (here).

To filter the messages, I think you can use the filter functions with the dide output streams :) (here)

Related