From 4c8d7725c5f5747a6a5250596f64a9f8297029c4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jose=20Luis=20L=C3=B3pez=20L=C3=B3pez?= Date: Fri, 25 Sep 2026 12:34:49 +0200 Subject: [PATCH] YARN-11993. PublicLocalizer thread exits permanently when pending.remove() returns null When pending.remove(completed) returned null, run() logged "Localized unknown resource" and returned. That terminated the Public Localizer thread for the life of the NodeManager, and run()'s finally block shut the download pool down on the way out. Every later public-resource request was then rejected; YARN-1800's handler turns each RejectedExecutionException into a ResourceFailedLocalizationEvent, so containers fail rather than hang, but no public resource can be localized on that node again until the NM is restarted. That is the state reported in YARN-9968, which fixed a different branch of this method. Skip the unattributable resource with continue instead. The branch is defensive rather than reachable today: addResource() wraps queue.submit(...) and pending.put(...) in one synchronized (pending) block so a future cannot complete and be dequeued before the map is updated, and pending is a synchronizedMap on the same monitor. It matters because YARN-11746 proposes shutting the NodeManager down when this thread exits, which would turn a single unknown resource into a node-wide outage. Hoisting the null check out of the try block also makes the assoc dereference in the ExecutionException handler provably safe, so SpotBugs no longer reports NP_NULL_ON_SOME_PATH_EXCEPTION. The exclude entry for it is therefore dropped rather than repaired; it had been inert since the revert of HADOOP-19668/19670 (#8568) restored run() without restoring the method name YARN-11912 changed to work(). The new test submits a download straight to the completion queue so the Future is never recorded in pending, then asserts the localizer neither dies nor shuts its pool down. It fails on the unfixed code with "Public Localizer exited after taking an unknown resource". Co-Authored-By: Claude Opus 5 --- .../dev-support/findbugs-exclude.xml | 6 -- .../ResourceLocalizationService.java | 10 +-- .../TestResourceLocalizationService.java | 69 +++++++++++++++++++ 3 files changed, 74 insertions(+), 11 deletions(-) diff --git a/hadoop-yarn-project/hadoop-yarn/dev-support/findbugs-exclude.xml b/hadoop-yarn-project/hadoop-yarn/dev-support/findbugs-exclude.xml index bf7440c61cb59c..0d82f357680c33 100644 --- a/hadoop-yarn-project/hadoop-yarn/dev-support/findbugs-exclude.xml +++ b/hadoop-yarn-project/hadoop-yarn/dev-support/findbugs-exclude.xml @@ -789,12 +789,6 @@ - - - - - - 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..324ac6953b89ce 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 @@ -982,12 +982,12 @@ public void run() { try { Future completed = queue.take(); LocalizerResourceRequestEvent assoc = pending.remove(completed); + if (null == assoc) { + LOG.error("Localized unknown resource to " + completed); + // TODO delete + continue; + } try { - if (null == assoc) { - LOG.error("Localized unknown resource to " + completed); - // TODO delete - return; - } Path local = completed.get(); LocalResourceRequest key = assoc.getResource().getRequest(); publicRsrc.handle(new ResourceLocalizedEvent(key, local, FileUtil 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..bca984718ef669 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 @@ -63,8 +63,11 @@ import java.util.Random; import java.util.Set; import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.CyclicBarrier; import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; import org.apache.hadoop.util.Sets; @@ -87,6 +90,7 @@ import org.apache.hadoop.io.DataOutputBuffer; import org.apache.hadoop.io.Text; import org.apache.hadoop.ipc.Server; +import org.apache.hadoop.test.GenericTestUtils; import org.apache.hadoop.security.Credentials; import org.apache.hadoop.security.token.Token; import org.apache.hadoop.security.token.TokenIdentifier; @@ -2650,6 +2654,71 @@ public void testParallelDownloadAttemptsForPublicResource() throws Exception { } + /** + * The Public Localizer must stay up when it dequeues a completed download it + * has no record of. Before YARN-11993 run() returned on that path, so the + * finally block shut the download pool down and every later public + * localization on the node was rejected until the NodeManager restarted. + */ + @Test + @Timeout(value = 30) + public void testPublicLocalizerSurvivesUnknownResource() throws Exception { + conf.setStrings(YarnConfiguration.NM_LOCAL_DIRS, + lfs.makeQualified(new Path(basedir, "0")).toString()); + + DrainDispatcher dispatcher = new DrainDispatcher(); + dispatcher.init(conf); + dispatcher.start(); + + // Nothing is localized here, so the dirs handler is never asked for a + // path; mocking it keeps the test off the disk. + LocalDirsHandlerService mockDirsHandler = + mock(LocalDirsHandlerService.class); + + ResourceLocalizationService service = + new ResourceLocalizationService(dispatcher, + mock(ContainerExecutor.class), mock(DeletionService.class), + mockDirsHandler, nmContext, metrics); + dispatcher.register(LocalizationEventType.class, service); + service.init(conf); + + PublicLocalizer publicLocalizer = service.getPublicLocalizer(); + try { + publicLocalizer.start(); + + // Submit straight to the completion queue so the Future is never + // recorded in pending. That is exactly the state in which + // pending.remove(completed) returns null. + final CountDownLatch downloaded = new CountDownLatch(1); + final Path unknown = new Path(basedir, "unknown"); + publicLocalizer.queue.submit(() -> { + downloaded.countDown(); + return unknown; + }); + assertTrue(downloaded.await(10, TimeUnit.SECONDS), + "public download never ran"); + assertEquals(0, publicLocalizer.pending.size()); + + // The localizer should log the unknown resource and carry on. If it + // exits instead, run()'s finally block shuts the download pool down. + try { + GenericTestUtils.waitFor(() -> !publicLocalizer.isAlive() + || publicLocalizer.threadPool.isShutdown(), 20, 5000); + fail("Public Localizer exited after taking an unknown resource"); + } catch (TimeoutException expected) { + // The localizer stayed up, which is what YARN-11993 fixed. + } + + assertTrue(publicLocalizer.isAlive(), "Public Localizer thread died"); + assertFalse(publicLocalizer.threadPool.isShutdown(), + "Public Localizer shut its download pool down"); + } finally { + publicLocalizer.interrupt(); + service.stop(); + dispatcher.stop(); + } + } + private boolean waitForPrivateDownloadToStart( ResourceLocalizationService service, String localizerId, int size, int maxWaitTime) {