When we restart a Flink job, the job needs to process millions of outstanding messages from several Kafka topics that have accumulated when the Flink job was down. Suddenly, as soon as the job starts, there will be a burst of input messages into the Flink job. The challenge we have is to process those messages on multiple topics in the order they occurred, i.e. process based on their event time. What are the best ways/solutions to handle this scenario. Is it a good idea to sort/partition the messages by a key and sort them based on the event time using a PriorityQueue and define a window of 1 minute or more and output the messages for the downstream processing. Or there any better solutions to address this? Is there a way to resolve this using watermarks?
The messages may have been created/written a few hours apart.