From 3ab8bc384d9c6afde97b6e85a8acffca39dea7ca Mon Sep 17 00:00:00 2001 From: Martin Vu <22mvu7@gmail.com> Date: Wed, 29 Jul 2026 23:50:53 -0700 Subject: [PATCH 1/4] Add timeout handling for region termination retries --- .../scheduling/RegionExecutionManager.scala | 28 +++++++++++-------- 1 file changed, 17 insertions(+), 11 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala index 93d085b6113..d81d2a281b3 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala @@ -172,7 +172,9 @@ class RegionExecutionManager( } private def terminateWorkers(regionExecution: RegionExecution) = { - // 1. Send EndWorkers to every worker + implicit val timer: Timer = new JavaTimer(true) + val killTimeout = TwitterDuration(30, TimeUnit.SECONDS) + // 1. Send EndWorkers with timeout val endWorkerRequests = regionExecution.getAllOperatorExecutions.flatMap { case (_, opExec) => @@ -183,9 +185,11 @@ class RegionExecutionManager( }.toSeq val endWorkerFuture: Future[Unit] = - Future.collect(endWorkerRequests).unit + Future.collect(endWorkerRequests) + .within(killTimeout) + .unit - // 2. Send GracefulStops only after 1 has finished + // 2. Send GracefulStops with timeout val gracefulStopRequests: Future[Unit] = endWorkerFuture.flatMap { _ => val gracefulStops = @@ -193,20 +197,16 @@ class RegionExecutionManager( case (_, opExec) => opExec.getWorkerIds.map { workerId => val actorRef = actorRefService.getActorRef(workerId) - // Remove the actorRef so that no other actors can find the worker and send messages. - actorRefService.removeActorRef(workerId) - // Restarted regions reuse actorId. Remove stale control channels so the - // coordinator does not reuse old control-message sequence numbers for new workers. - asyncRPCClient.inputGateway.removeControlChannel(workerId) - asyncRPCClient.outputGateway.removeControlChannel(workerId) gracefulStop(actorRef, ScalaDuration(5, TimeUnit.SECONDS)).asTwitter() } }.toSeq - Future.collect(gracefulStops).unit + Future.collect(gracefulStops) + .within(killTimeout) + .unit } - // 3. Log whether the kills were successful + // 3. Cleanup only after all gracefulStops succeed gracefulStopRequests.transform { case Return(_) => logger.debug(s"Region ${region.id.id} successfully terminated.") @@ -214,6 +214,12 @@ class RegionExecutionManager( case (_, opExec) => opExec.getWorkerIds.foreach { workerId => opExec.getWorkerExecution(workerId).forceTerminate() + // Remove the actorRef so that no other actors can find the worker and send messages. + actorRefService.removeActorRef(workerId) + // Restarted regions reuse actorId. Remove stale control channels so the + // coordinator does not reuse old control-message sequence numbers for new workers. + asyncRPCClient.inputGateway.removeControlChannel(workerId) + asyncRPCClient.outputGateway.removeControlChannel(workerId) } } Future.Unit // propagate success From a6c8b80acf479f9a654ca28b56aab1c66714ed72 Mon Sep 17 00:00:00 2001 From: Martin Vu <22mvu7@gmail.com> Date: Thu, 30 Jul 2026 00:14:31 -0700 Subject: [PATCH 2/4] style: apply scalafmt formatting to RegionExecutionManager --- .../architecture/scheduling/RegionExecutionManager.scala | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala index 47b38a010f0..a01ebc62325 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala @@ -198,7 +198,8 @@ class RegionExecutionManager( }.toSeq val endWorkerFuture: Future[Unit] = - Future.collect(endWorkerRequests) + Future + .collect(endWorkerRequests) .within(killTimeout) .unit @@ -214,7 +215,8 @@ class RegionExecutionManager( } }.toSeq - Future.collect(gracefulStops) + Future + .collect(gracefulStops) .within(killTimeout) .unit } From bdb3ad34d622c18acf9531511643ca7a6729b336 Mon Sep 17 00:00:00 2001 From: Martin Vu <22mvu7@gmail.com> Date: Thu, 30 Jul 2026 01:33:02 -0700 Subject: [PATCH 3/4] fix(amber): shorten EndWorker termination timeout --- .../architecture/scheduling/RegionExecutionManager.scala | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala index a01ebc62325..546ab4fb91e 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala @@ -186,7 +186,7 @@ class RegionExecutionManager( private def terminateWorkers(regionExecution: RegionExecution) = { implicit val timer: Timer = new JavaTimer(true) - val killTimeout = TwitterDuration(30, TimeUnit.SECONDS) + val killTimeout = com.twitter.util.Duration.fromMilliseconds(5000) // 1. Send EndWorkers with timeout val endWorkerRequests = regionExecution.getAllOperatorExecutions.flatMap { @@ -229,7 +229,7 @@ class RegionExecutionManager( case (_, opExec) => opExec.getWorkerIds.foreach { workerId => opExec.getWorkerExecution(workerId).forceTerminate() - // Remove the actorRef so that no other actors can find the worker and send messages. + // Remove the actorRef after successful termination so other actors cannot reach the worker. actorRefService.removeActorRef(workerId) // Restarted regions reuse actorId. Remove stale control channels so the // coordinator does not reuse old control-message sequence numbers for new workers. From 010367bd3be83bbd04f4249101298d6c18f64617 Mon Sep 17 00:00:00 2001 From: Martin Vu <22mvu7@gmail.com> Date: Thu, 30 Jul 2026 01:47:50 -0700 Subject: [PATCH 4/4] fix(amber): shorten EndWorker termination timeout --- .../engine/architecture/scheduling/RegionExecutionManager.scala | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala index 546ab4fb91e..74849cf0286 100644 --- a/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala +++ b/amber/src/main/scala/org/apache/texera/amber/engine/architecture/scheduling/RegionExecutionManager.scala @@ -186,7 +186,7 @@ class RegionExecutionManager( private def terminateWorkers(regionExecution: RegionExecution) = { implicit val timer: Timer = new JavaTimer(true) - val killTimeout = com.twitter.util.Duration.fromMilliseconds(5000) + val killTimeout = com.twitter.util.Duration.fromMilliseconds(1000) // 1. Send EndWorkers with timeout val endWorkerRequests = regionExecution.getAllOperatorExecutions.flatMap {