From 20d9107e885321f25fc43328ec1ccf5b8a667fa4 Mon Sep 17 00:00:00 2001 From: hyunw9 Date: Tue, 28 Jul 2026 23:12:42 +0900 Subject: [PATCH 1/2] [ZEPPELIN-6574] Add read-only REST API for interpreter process status --- .../interpreter/InterpreterProcessStatus.java | 97 +++++++++++++++++++ .../InterpreterSettingManager.java | 12 +++ .../remote/RemoteInterpreterProcess.java | 9 ++ .../zeppelin/rest/InterpreterRestApi.java | 11 +++ .../InterpreterSettingManagerTest.java | 24 +++++ .../zeppelin/rest/InterpreterRestApiTest.java | 11 +++ 6 files changed, 164 insertions(+) create mode 100644 zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java new file mode 100644 index 00000000000..4a6a1f9d634 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java @@ -0,0 +1,97 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.interpreter; + +import org.apache.zeppelin.interpreter.remote.RemoteInterpreterProcess; + +/** + * Point-in-time status snapshot of a single interpreter process as seen by the Zeppelin server. + * Built purely from in-memory server state without contacting the process, so {@code started} + * reflects whether a process handle exists, not whether the process is currently reachable. + * Reachability is intentionally out of scope here to keep the read path non-blocking. + */ +public class InterpreterProcessStatus { + private final String settingId; + private final String settingName; + private final String groupId; + private final int numSessions; + private final boolean started; + private String host; + private int port = -1; + private String startTime; + private long uptimeSeconds; + private String errorMessage; + + public InterpreterProcessStatus(ManagedInterpreterGroup group) { + InterpreterSetting setting = group.getInterpreterSetting(); + this.settingId = setting.getId(); + this.settingName = setting.getName(); + this.groupId = group.getId(); + this.numSessions = group.getSessionNum(); + // Read the handle once: another thread may close the group concurrently. + RemoteInterpreterProcess process = group.getInterpreterProcess(); + this.started = process != null; + if (started) { + this.host = process.getHost(); + this.port = process.getPort(); + this.startTime = process.getStartTime(); + this.uptimeSeconds = (System.currentTimeMillis() - process.getStartTimeMs()) / 1000; + this.errorMessage = process.getErrorMessage(); + } + } + + public String getSettingId() { + return settingId; + } + + public String getSettingName() { + return settingName; + } + + public String getGroupId() { + return groupId; + } + + public int getNumSessions() { + return numSessions; + } + + public boolean isStarted() { + return started; + } + + public String getHost() { + return host; + } + + public int getPort() { + return port; + } + + public String getStartTime() { + return startTime; + } + + public long getUptimeSeconds() { + return uptimeSeconds; + } + + public String getErrorMessage() { + return errorMessage; + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java index 08d11629ba8..35f3d8f91ae 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java @@ -701,6 +701,18 @@ public List getAllInterpreterGroup() { return interpreterGroups; } + /** + * Snapshot the status of every running interpreter group. Uses in-memory state only (no remote + * probe) so a stuck interpreter cannot block this call. + */ + public List getInterpreterProcessStatuses() { + List statuses = new ArrayList<>(); + for (ManagedInterpreterGroup group : getAllInterpreterGroup()) { + statuses.add(new InterpreterProcessStatus(group)); + } + return statuses; + } + // TODO(zjffdu) Current approach is not optimized. we have to iterate all interpreter settings. public void removeInterpreterGroup(String intpGroupId) { for (InterpreterSetting interpreterSetting : interpreterSettings.values()) { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java index 95802a64fe7..977f4bde804 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java @@ -43,6 +43,7 @@ public abstract class RemoteInterpreterProcess implements InterpreterClient, Aut protected int intpEventServerPort; private PooledRemoteClient remoteClient; private String startTime; + private final long startTimeMs; public RemoteInterpreterProcess(int connectTimeout, int connectionPoolSize, @@ -52,6 +53,7 @@ public RemoteInterpreterProcess(int connectTimeout, this.intpEventServerHost = intpEventServerHost; this.intpEventServerPort = intpEventServerPort; this.startTime = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date()); + this.startTimeMs = System.currentTimeMillis(); this.remoteClient = new PooledRemoteClient<>(() -> { TSocket transport = new TSocket(getHost(), getPort()); try { @@ -72,6 +74,13 @@ public String getStartTime() { return startTime; } + /** + * Epoch millis captured at construction, used to compute uptime without a remote call. + */ + public long getStartTimeMs() { + return startTimeMs; + } + @Override public void close() { if (remoteClient != null) { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java index 3b9d754e919..163fad96145 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java @@ -102,6 +102,17 @@ public Response listSettings() { return new JsonResponse<>(Status.OK, "", interpreterSettingManager.get()).build(); } + /** + * List the runtime status of all running interpreter processes. + */ + @GET + @Path("status") + @ZeppelinApi + public Response getInterpreterProcessStatus() { + return new JsonResponse<>(Status.OK, "", + interpreterSettingManager.getInterpreterProcessStatuses()).build(); + } + /** * Get a setting. */ diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java index 95e126fedff..39f695c012e 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java @@ -40,7 +40,9 @@ import java.util.Map; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; 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.junit.jupiter.api.Assertions.fail; import static org.mockito.Mockito.mock; @@ -262,6 +264,28 @@ void testRestartShared() throws InterpreterException { assertEquals(0, interpreterSetting.getAllInterpreterGroups().size()); } + @Test + void testGetInterpreterProcessStatuses() throws InterpreterException { + // no interpreter group has been created yet + assertTrue(interpreterSettingManager.getInterpreterProcessStatuses().isEmpty()); + + InterpreterSetting interpreterSetting = interpreterSettingManager.getByName("test"); + interpreterSetting.getOption().setPerUser("shared"); + interpreterSetting.getOption().setPerNote("shared"); + interpreterSetting.getOrCreateSession("user1", note1Id); + + List statuses = + interpreterSettingManager.getInterpreterProcessStatuses(); + assertEquals(1, statuses.size()); + InterpreterProcessStatus status = statuses.get(0); + assertEquals("test", status.getSettingName()); + assertEquals(1, status.getNumSessions()); + // process starts lazily on first interpret, so it is not started at this point + assertFalse(status.isStarted()); + assertNull(status.getHost()); + assertEquals(-1, status.getPort()); + } + @Test void testRestartPerUserIsolated() throws InterpreterException { InterpreterSetting interpreterSetting = interpreterSettingManager.getByName("test"); diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java index 19435b32112..66d266b703d 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java @@ -106,6 +106,17 @@ void getSettings() throws IOException { get.close(); } + @Test + void testGetInterpreterProcessStatus() throws IOException { + // when + CloseableHttpResponse get = httpGet("/interpreter/status"); + // then + assertThat(get, isAllowed()); + JsonArray body = getArrayBodyFieldFromResponse(EntityUtils.toString(get.getEntity(), StandardCharsets.UTF_8)); + assertNotNull(body); + get.close(); + } + @Test void testGetNonExistInterpreterSetting() throws IOException { // when From 679ed0823601ab3cfc098530269c826c14fc9db2 Mon Sep 17 00:00:00 2001 From: hyunw9 Date: Sun, 9 Aug 2026 17:10:12 +0900 Subject: [PATCH 2/2] [ZEPPELIN-6574] Remove unnecessary comments --- .../zeppelin/interpreter/InterpreterProcessStatus.java | 1 - .../zeppelin/interpreter/InterpreterSettingManager.java | 3 +-- .../interpreter/remote/RemoteInterpreterProcess.java | 5 +---- 3 files changed, 2 insertions(+), 7 deletions(-) diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java index 4a6a1f9d634..2330df6bd4c 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java @@ -43,7 +43,6 @@ public InterpreterProcessStatus(ManagedInterpreterGroup group) { this.settingName = setting.getName(); this.groupId = group.getId(); this.numSessions = group.getSessionNum(); - // Read the handle once: another thread may close the group concurrently. RemoteInterpreterProcess process = group.getInterpreterProcess(); this.started = process != null; if (started) { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java index 35f3d8f91ae..bf8ae8ffcf5 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java @@ -702,8 +702,7 @@ public List getAllInterpreterGroup() { } /** - * Snapshot the status of every running interpreter group. Uses in-memory state only (no remote - * probe) so a stuck interpreter cannot block this call. + * Snapshot the status of every running interpreter group. Uses in-memory state only */ public List getInterpreterProcessStatuses() { List statuses = new ArrayList<>(); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java index 977f4bde804..4c8c3b63285 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java @@ -73,10 +73,7 @@ public int getConnectTimeout() { public String getStartTime() { return startTime; } - - /** - * Epoch millis captured at construction, used to compute uptime without a remote call. - */ + public long getStartTimeMs() { return startTimeMs; }