From 3067fdee3eeb1275ce3c8dffb58711f5739dcacd Mon Sep 17 00:00:00 2001 From: JangAyeon Date: Fri, 25 Sep 2026 11:44:54 +0900 Subject: [PATCH] [ZEPPELIN-6693] Drive NotebookServer heartbeat scheduler shutdown from ZeppelinServer lifecycle The heartbeat scheduler was stopped by its own JVM shutdown hook, which never ran when the server was closed inside a running JVM (e.g. in tests), so schedulers and hooks piled up across start/stop cycles. ZeppelinServer#shutdown now calls stopHeartbeatScheduler() after Jetty stops, covering both in-JVM close and JVM exit. The dedicated hook is removed. Add a test that checks the scheduler is stopped via MiniZeppelinServer.shutDown(). --- .../zeppelin/server/ZeppelinServer.java | 2 + .../zeppelin/socket/NotebookServer.java | 18 ++------ .../socket/NotebookServerHeartbeatTest.java | 46 +++++++++++++++++++ 3 files changed, 51 insertions(+), 15 deletions(-) diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java index 6dece8a13d3..480422b4a5b 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/server/ZeppelinServer.java @@ -359,6 +359,8 @@ public void shutdown(int exitCode) { jettyWebServer.stop(); } if (sharedServiceLocator != null) { + // Stop after Jetty so no new connection can restart the heartbeat scheduler. + sharedServiceLocator.getService(NotebookServer.class).stopHeartbeatScheduler(); if (!zConf.isRecoveryEnabled()) { sharedServiceLocator.getService(InterpreterSettingManager.class).close(); } diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java index 4d31a06558c..741a8348242 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java @@ -150,7 +150,6 @@ String getKey() { // Package-private (not private) so NotebookServerHeartbeatTest can observe scheduler // lifecycle without exposing it as part of the public API. ScheduledExecutorService heartbeatScheduler; - private Thread heartbeatShutdownHook; private boolean heartbeatInitialized; // TODO(jl): This will be removed by handling session directly @@ -287,29 +286,18 @@ synchronized void startHeartbeatScheduler() { }); heartbeatScheduler.scheduleAtFixedRate( this::sendHeartbeat, intervalMs, intervalMs, TimeUnit.MILLISECONDS); - heartbeatShutdownHook = new Thread(this::stopHeartbeatScheduler); - Runtime.getRuntime().addShutdownHook(heartbeatShutdownHook); LOGGER.info("Started websocket heartbeat scheduler with interval {} ms", intervalMs); } /** - * Stops the websocket heartbeat scheduler, if running, and deregisters its shutdown hook so - * repeated start/stop cycles do not accumulate hooks. Safe to call multiple times and safe - * to call when the scheduler was never started. + * Stops the websocket heartbeat scheduler, if running. Safe to call multiple times and safe + * to call when the scheduler was never started; a later connection starts it again. */ - synchronized void stopHeartbeatScheduler() { + public synchronized void stopHeartbeatScheduler() { if (heartbeatScheduler != null) { heartbeatScheduler.shutdownNow(); heartbeatScheduler = null; } - if (heartbeatShutdownHook != null && Thread.currentThread() != heartbeatShutdownHook) { - try { - Runtime.getRuntime().removeShutdownHook(heartbeatShutdownHook); - } catch (IllegalStateException e) { - // JVM is already shutting down; the hook will simply run (as a harmless no-op). - } - heartbeatShutdownHook = null; - } heartbeatInitialized = false; } diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java index 6f53668261b..f932b1fa2dd 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java @@ -19,11 +19,14 @@ import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import java.util.concurrent.ScheduledExecutorService; +import org.apache.zeppelin.MiniZeppelinServer; import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.notebook.AuthorizationService; import org.junit.jupiter.api.AfterEach; @@ -106,4 +109,47 @@ void startHeartbeatSchedulerDoesNotStartWhenIntervalIsNegative() { assertNull(server.heartbeatScheduler); } + + @Test + void stopHeartbeatSchedulerAllowsRepeatedStartStopCycles() { + NotebookServer server = buildNotebookServer(50L); + + for (int i = 0; i < 3; i++) { + server.startHeartbeatScheduler(); + ScheduledExecutorService scheduler = server.heartbeatScheduler; + assertNotNull(scheduler); + + server.stopHeartbeatScheduler(); + + assertTrue(scheduler.isShutdown()); + assertNull(server.heartbeatScheduler); + } + } + + @Test + void stopHeartbeatSchedulerIsSafeWhenNeverStarted() { + NotebookServer server = buildNotebookServer(50L); + + assertDoesNotThrow(server::stopHeartbeatScheduler); + assertDoesNotThrow(server::stopHeartbeatScheduler); + } + + @Test + void zeppelinServerShutdownStopsHeartbeatScheduler() throws Exception { + MiniZeppelinServer zepServer = + new MiniZeppelinServer(NotebookServerHeartbeatTest.class.getSimpleName()); + try { + zepServer.start(); + NotebookServer server = zepServer.getService(NotebookServer.class); + server.startHeartbeatScheduler(); + ScheduledExecutorService scheduler = server.heartbeatScheduler; + assertNotNull(scheduler); + + zepServer.shutDown(); + + assertTrue(scheduler.isShutdown()); + } finally { + zepServer.destroy(); + } + } }