Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,111 @@
/*
* 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.
*
* <p>Every value below must be readable without leaving the JVM. In particular do not call
* {@code isRunning()}, {@code isAlive()} or {@code getErrorMessage()} on the process from here:
* those are failure-path diagnostics that contact the container runtime on some launchers, so
* calling them would let a slow or unreachable runtime block this endpoint.
*
* <p>{@code started} and {@code launching} are read one after the other rather than under the
* lock that guards a launch, so this is a best-effort view of a group that is starting up: the
* pair can straddle the moment a launch finishes. What it does buy is that the window in which
* a handle carries no {@code host} or {@code port} yet is reported as such instead of looking
* like a fully started process.
*/
public class InterpreterProcessStatus {
private final String settingId;
private final String settingName;
private final String groupId;
private final int numSessions;
private final boolean launching;
private final boolean started;
private String host;
private int port = -1;
private String startTime;
private long attachedForSeconds;

public InterpreterProcessStatus(ManagedInterpreterGroup group) {
InterpreterSetting setting = group.getInterpreterSetting();
this.settingId = setting.getId();
this.settingName = setting.getName();
this.groupId = group.getId();
this.numSessions = group.getSessionNum();
this.launching = group.isLaunchingInterpreterProcess();
RemoteInterpreterProcess process = group.getInterpreterProcess();
this.started = process != null;
if (started) {
Comment on lines +59 to +60

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ManagedInterpreterGroup.getOrCreateInterpreterProcess() assigns remoteInterpreterProcess = createInterpreterProcess(...) before start(), and the reader does not take interpreterProcessCreationLock. A call during a launch can therefore observe a handle whose host and port are still at their initial values (null / -1 at RemoteInterpreterManagedProcess:35-36).

A single started boolean does not let a consumer tell that window apart from a fully started process, so you may want a second field. Conveniently ManagedInterpreterGroup.isLaunchingInterpreterProcess() already exists, or the state could be derived from whether port has been filled in. The latter looks more robust for separating "handle created / awaiting registration / registered" at no extra cost, though you may prefer the simplicity of one boolean, so I will leave it as a matter of taste.

Combined with the uptime in the note above, this would also make "stuck awaiting registration" visible on its own, which catches a fair amount without any remote probe.

this.host = process.getHost();
this.port = process.getPort();
Comment on lines +61 to +62

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Most of the values this snapshot reads are non-volatile and are written on a different thread from the one reading them:

  • ManagedInterpreterGroup.remoteInterpreterProcess: written inside a synchronized block, read without the lock
  • RemoteInterpreterManagedProcess.host / port: written by the Thrift registration callback, read by the REST thread
  • errorMessage: written by the YarnAppMonitor scheduler thread or the K8s path, read by the REST thread

All of this predates the PR, so it is not something this change introduced. I mention it because this API is the first place that state becomes a documented contract, so reporting a stale value now has a visible consequence. It also feeds directly into the phase decision if you take the port-based approach above.

A few volatile modifiers look close to free here, but if that feels out of scope it seems fine as follow-up. Your call.

this.startTime = process.getStartTime();
this.attachedForSeconds = (System.currentTimeMillis() - process.getStartTimeMs()) / 1000;
}
}

public String getSettingId() {
return settingId;
}

public String getSettingName() {
return settingName;
}

public String getGroupId() {
return groupId;
}

public int getNumSessions() {
return numSessions;
}

/**
* @return whether a process is currently being launched for this group, in which case
* {@code host} and {@code port} may not be filled in yet even when {@code started}
*/
public boolean isLaunching() {
return launching;
}

public boolean isStarted() {
return started;
}

public String getHost() {
return host;
}

public int getPort() {
return port;
}

public String getStartTime() {
return startTime;
}

public long getAttachedForSeconds() {
return attachedForSeconds;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -711,6 +711,17 @@ public List<ManagedInterpreterGroup> getAllInterpreterGroup() {
return interpreterGroups;
}

/**
* Snapshot the status of every running interpreter group. Uses in-memory state only
*/
public List<InterpreterProcessStatus> getInterpreterProcessStatuses() {
List<InterpreterProcessStatus> 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()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ public class ManagedInterpreterGroup extends InterpreterGroup {
private static final Logger LOGGER = LoggerFactory.getLogger(ManagedInterpreterGroup.class);

private InterpreterSetting interpreterSetting;
private RemoteInterpreterProcess remoteInterpreterProcess; // attached remote interpreter process
private volatile RemoteInterpreterProcess remoteInterpreterProcess;
private Object interpreterProcessCreationLock = new Object();
private final ZeppelinConfiguration zConf;
private volatile long lastUsedTimeInMillis = System.currentTimeMillis();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,14 +32,14 @@ public abstract class RemoteInterpreterManagedProcess extends RemoteInterpreterP

private final String interpreterPortRange;

private String host = null;
private int port = -1;
private volatile String host = null;
private volatile int port = -1;
private final String interpreterDir;
private final String localRepoDir;
private final String interpreterSettingName;
private final String interpreterGroupId;
private final boolean isUserImpersonated;
private String errorMessage;
private volatile String errorMessage;

private Map<String, String> env;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,8 @@ public abstract class RemoteInterpreterProcess implements InterpreterClient, Aut
protected String intpEventServerHost;
protected int intpEventServerPort;
private PooledRemoteClient<Client> remoteClient;
private String startTime;
private final long startTimeMs;
private final String startTime;

public RemoteInterpreterProcess(int connectTimeout,
int connectionPoolSize,
Expand All @@ -54,7 +55,8 @@ public RemoteInterpreterProcess(int connectTimeout,
this.connectTimeout = 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.startTime = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date(startTimeMs));
this.remoteClient = new PooledRemoteClient<>(() -> {
TSocket transport = new TSocket(getHost(), getPort());
try {
Expand All @@ -71,9 +73,21 @@ public int getConnectTimeout() {
return connectTimeout;
}

/**
* When the server created this object, formatted for display. This is not necessarily when the
* interpreter itself started: {@link RemoteInterpreterRunningProcess} is constructed fresh when
* the server recovers a process that outlived it, and when it attaches to an interpreter that
* was already running, so on those paths the stamp is the moment of attachment.
*
* @return the creation instant of this object as {@code yyyy-MM-dd HH:mm:ss}
*/
public String getStartTime() {
return startTime;
}

public long getStartTimeMs() {
return startTimeMs;
}

@Override
public void close() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -262,6 +264,29 @@ 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<InterpreterProcessStatus> 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());
assertFalse(status.isLaunching());
assertNull(status.getHost());
assertEquals(-1, status.getPort());
}

@Test
void testRestartPerUserIsolated() throws InterpreterException {
InterpreterSetting interpreterSetting = interpreterSettingManager.getByName("test");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading