Hi Sandeep! While auto scaling jobs in Flink still isn’t possible, in Flink 1.2 you will be able to rescale jobs by stopping and restarting. This works by taking a savepoint of the job before stopping the job, and then redeploy the job with a higher / lower parallelism using the savepoint. Upon restarting the job, your states will be redistributed across the new operators.
Changing operator / job parallelism on the fly while running is still on the future roadmap. Cheers, Gordon On February 2, 2017 at 8:39:39 AM, Meghashyam Sandeep V (vr1meghash...@gmail.com) wrote: Hi Guys, I currently run flink 1.1.4 streaming jobs in EMR in AWS with yarn. I understand that EMR supports auto scaling but Flink doesn't. Is there a plan for this support in 1.2. Thanks, Sandeep