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) {