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 ([hidden email]) 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