Hello All, My name is Tarek, I have been working with Flink for over a year now, we are facing an issue with one of our flink job that have many joins, the original table we are starting with is about 430 millions but joining with other tables that are between 40-50 millions down to couple of 1000s records then eventually the sink flink will output the data to Postgres database. I know there will be alot of questions about some details needed here but this is at high level
We are using Flink 1.16 hosted in EMR cluster using rocksdb a state backend with about 30 tasks nodes for i4i.xlagre instances, the job I raised above have is set to 12 parallelism
The issue that we are facing that slowness is between the joins, during the checkpoints. If someone has try Flink with complex query joining and can share some of the settings used that improve their performance. I will be thankful. Also if someone wanted to dig deep into this with me, I will be grateful. gratitude thank you