From adf120e0e59a9c90f84e8e558d5934589b4845d8 Mon Sep 17 00:00:00 2001 From: Yubi Lee Date: Mon, 17 Aug 2026 15:18:52 +0900 Subject: [PATCH 1/2] YARN-11856. DOWNLOADING resources unlock and cleanup is interrupted when killing a container that is localizing. Contributed by zheng-weihao and Yubi Lee. --- .../ResourceLocalizationService.java | 8 ++- .../TestResourceLocalizationService.java | 54 +++++++++++++++++++ 2 files changed, 60 insertions(+), 2 deletions(-) diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java index a7f0722e66f8e0..cfb9919ab14b2a 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java @@ -1294,8 +1294,12 @@ public void run() { // On error, report failure to Container and signal ABORT // Notify resource of failed localization ContainerId cId = context.getContainerId(); - dispatcher.getEventHandler().handle(new ContainerResourceFailedEvent( - cId, null, exception.getMessage())); + try { + dispatcher.getEventHandler().handle(new ContainerResourceFailedEvent( + cId, null, exception.getMessage())); + } catch (Exception e) { + LOG.info("Failed to send container resource failed event for " + cId.toString(), e); + } } List paths = new ArrayList(); for (LocalizerResourceRequestEvent event : scheduled.values()) { diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java index e360a73f196b5c..f2c76c053c40a1 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java @@ -70,6 +70,7 @@ import org.apache.hadoop.util.Sets; import org.apache.hadoop.yarn.server.nodemanager.containermanager.container.ContainerState; import org.apache.hadoop.yarn.server.nodemanager.containermanager.deletion.task.FileDeletionMatcher; +import org.apache.hadoop.yarn.server.nodemanager.containermanager.deletion.task.FileDeletionTask; import org.apache.hadoop.yarn.server.nodemanager.executor.LocalizerStartContext; import org.apache.commons.io.FileUtils; import org.apache.hadoop.conf.Configuration; @@ -102,9 +103,12 @@ import org.apache.hadoop.yarn.api.records.URL; import org.apache.hadoop.yarn.conf.YarnConfiguration; import org.apache.hadoop.yarn.event.AsyncDispatcher; +import org.apache.hadoop.yarn.event.Dispatcher; import org.apache.hadoop.yarn.event.DrainDispatcher; +import org.apache.hadoop.yarn.event.Event; import org.apache.hadoop.yarn.event.EventHandler; import org.apache.hadoop.yarn.exceptions.YarnException; +import org.apache.hadoop.yarn.exceptions.YarnRuntimeException; import org.apache.hadoop.yarn.server.nodemanager.ContainerExecutor; import org.apache.hadoop.yarn.server.nodemanager.DefaultContainerExecutor; import org.apache.hadoop.yarn.server.nodemanager.DeletionService; @@ -818,6 +822,56 @@ public void testLocalizerRunnerException() throws Exception { } } + @Test + @Timeout(value = 10) + @SuppressWarnings("unchecked") // mocked generics + public void testDownloadingResourcesCleanedUpWhenDispatchFails() + throws Exception { + Dispatcher dispatcher = mock(Dispatcher.class); + EventHandler eventHandler = mock(EventHandler.class); + when(dispatcher.getEventHandler()).thenReturn(eventHandler); + // Simulate the localizer thread being interrupted by a container kill: + // dispatching the failure event throws instead of completing. + Mockito.doThrow(new YarnRuntimeException(new InterruptedException())) + .when(eventHandler).handle(isA(ContainerResourceFailedEvent.class)); + + ContainerExecutor exec = mock(ContainerExecutor.class); + DeletionService delService = mock(DeletionService.class); + LocalDirsHandlerService dirsHandlerSpy = spy(new LocalDirsHandlerService()); + dirsHandlerSpy.init(conf); + // Fail localization so LocalizerRunner.run() takes the error path. + Mockito.doThrow(new IOException("Simulated disk failure")) + .when(dirsHandlerSpy).getLocalPathForWrite(isA(String.class)); + + ResourceLocalizationService rls = + new ResourceLocalizationService(dispatcher, exec, delService, + dirsHandlerSpy, nmContext, metrics); + + final ApplicationId appId = + BuilderUtils.newApplicationId(314159265358979L, 3); + final Container c = getMockContainer(appId, 42, "user0"); + LocalizerRunner runner = rls.new LocalizerRunner( + new LocalizerContext("user0", c.getContainerId(), null), + c.getContainerId().toString()); + + // A resource that was in DOWNLOADING state when the localizer died. + LocalizedResource rsrc = mock(LocalizedResource.class); + when(rsrc.getLocalPath()).thenReturn( + new Path("/local/usercache/user0/filecache/10/foo.jar")); + LocalizerResourceRequestEvent scheduledEvent = + mock(LocalizerResourceRequestEvent.class); + when(scheduledEvent.getResource()).thenReturn(rsrc); + runner.scheduled.put(mock(LocalResourceRequest.class), scheduledEvent); + + // Must not propagate the dispatch failure, and must still unlock the + // DOWNLOADING resource and schedule the deletion tasks. + runner.run(); + + verify(rsrc).unlock(); + verify(delService, Mockito.atLeastOnce()) + .delete(isA(FileDeletionTask.class)); + } + @Test @Timeout(value = 10) @SuppressWarnings({"unchecked", "methodlength"}) // mocked generics From 8600dbb6f41c7aec93805a1888a631023988b00f Mon Sep 17 00:00:00 2001 From: Yubi Lee Date: Sat, 22 Aug 2026 19:16:35 +0900 Subject: [PATCH 2/2] YARN-11856. Address review comments: restore interrupt status, WARN level logging, stronger test assertions. - Restore the thread interrupt status when the swallowed dispatch failure was caused by an InterruptedException. - Log the dispatch failure at WARN with parameterized logging. - Assert the exact FileDeletionTasks (localization dir + _tmp dir, nmPrivate token file) and the restored interrupt status in the test. --- .../ResourceLocalizationService.java | 7 +++++- .../TestResourceLocalizationService.java | 24 +++++++++++++++++-- 2 files changed, 28 insertions(+), 3 deletions(-) diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java index cfb9919ab14b2a..b5c5b807d460a5 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/main/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/ResourceLocalizationService.java @@ -1298,7 +1298,12 @@ public void run() { dispatcher.getEventHandler().handle(new ContainerResourceFailedEvent( cId, null, exception.getMessage())); } catch (Exception e) { - LOG.info("Failed to send container resource failed event for " + cId.toString(), e); + LOG.warn("Failed to send container resource failed event for {}", + cId, e); + if (e instanceof InterruptedException + || e.getCause() instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } } } List paths = new ArrayList(); diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java index f2c76c053c40a1..df8f97ccc1950c 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-nodemanager/src/test/java/org/apache/hadoop/yarn/server/nodemanager/containermanager/localizer/TestResourceLocalizationService.java @@ -867,9 +867,29 @@ public void testDownloadingResourcesCleanedUpWhenDispatchFails() // DOWNLOADING resource and schedule the deletion tasks. runner.run(); + // The interrupt status must be restored after the dispatch failure was + // swallowed. Thread.interrupted() also clears it for the test thread. + assertTrue(Thread.interrupted()); + verify(rsrc).unlock(); - verify(delService, Mockito.atLeastOnce()) - .delete(isA(FileDeletionTask.class)); + ArgumentCaptor captor = + ArgumentCaptor.forClass(FileDeletionTask.class); + verify(delService, times(2)).delete(captor.capture()); + List tasks = captor.getAllValues(); + // Localization dir and _tmp download dir of the DOWNLOADING resource. + FileDeletionTask rsrcTask = tasks.get(0); + assertEquals("user0", rsrcTask.getUser()); + assertNull(rsrcTask.getSubDir()); + assertEquals(Arrays.asList( + new Path("/local/usercache/user0/filecache/10"), + new Path("/local/usercache/user0/filecache/10_tmp")), + rsrcTask.getBaseDirs()); + // nmPrivate token file; the path is null here because localization + // failed before it was resolved. + FileDeletionTask tokenTask = tasks.get(1); + assertNull(tokenTask.getUser()); + assertNull(tokenTask.getSubDir()); + assertNull(tokenTask.getBaseDirs()); } @Test