Skip to content
Open
Prev Previous commit
[ZEPPELIN-6709] Clarify output ordering and strengthen boundary tests
  • Loading branch information
miinhho committed Sep 27, 2026
commit f164e310be64c14e39dd03fe668d25808348df50
5 changes: 3 additions & 2 deletions conf/zeppelin-site.xml.template
Original file line number Diff line number Diff line change
Expand Up @@ -443,8 +443,9 @@
<name>zeppelin.interpreter.output.worker.count</name>
Comment thread
voidmatcha marked this conversation as resolved.
<value>4</value>
<description>Number of workers in each server pool: one pool delivers paragraph output, and the
other saves output checkpoints. Output for a note stays in order, while different notes can be
processed concurrently. Checkpoint saves do not occupy output workers.</description>
other saves output checkpoints. Each paragraph's output stays in order, and updates or checkpoints
never overtake earlier appends. Different notes can be processed concurrently. Checkpoint saves
do not occupy output workers.</description>
</property>

<property>
Expand Down
2 changes: 1 addition & 1 deletion docs/setup/operation/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -374,7 +374,7 @@ Sources descending by priority:
<td><h6 class="properties">ZEPPELIN_INTERPRETER_OUTPUT_WORKER_COUNT</h6></td>
<td><h6 class="properties">zeppelin.interpreter.output.worker.count</h6></td>
<td>4</td>
<td>Number of workers in each server pool: one pool delivers paragraph output, and the other saves output checkpoints. Output for a note stays in order, while different notes can be processed concurrently. Checkpoint saves do not occupy output workers.</td>
<td>Number of workers in each server pool: one pool delivers paragraph output, and the other saves output checkpoints. Each paragraph's output stays in order, and updates or checkpoints never overtake earlier appends. Different notes can be processed concurrently. Checkpoint saves do not occupy output workers.</td>
</tr>
<tr>
<td><h6 class="properties">ZEPPELIN_INTERPRETER_OUTPUT_EVENTS_PER_BATCH</h6></td>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
*/
package org.apache.zeppelin.interpreter;

import static org.awaitility.Awaitility.await;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
Expand Down Expand Up @@ -45,7 +46,6 @@
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.zeppelin.conf.ZeppelinConfiguration;
Expand Down Expand Up @@ -174,26 +174,67 @@ void sameNoteBoundaryWaitsForInFlightAppend(String operation) throws Exception {
RemoteInterpreterEventServer server = serverWithListener(listener);
CountDownLatch entered = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
CountDownLatch started = new CountDownLatch(1);
AtomicReference<Thread> callerThread = new AtomicReference<>();
AtomicBoolean appendCallbackActive = new AtomicBoolean();
AtomicBoolean appendCompleted = new AtomicBoolean();
AtomicBoolean boundaryOverlappedAppend = new AtomicBoolean();
AtomicBoolean rpcReturnedEarly = new AtomicBoolean();

doAnswer(invocation -> {
appendCallbackActive.set(true);
entered.countDown();
assertTrue(release.await(5, TimeUnit.SECONDS));
try {
assertTrue(release.await(5, TimeUnit.SECONDS));
} finally {
appendCallbackActive.set(false);
appendCompleted.set(true);
}
return null;
}).when(listener).onParagraphOutputAppend("note", "first", 0, null, "old");
doAnswer(call -> {
if (appendCallbackActive.get() || !appendCompleted.get()) {
boundaryOverlappedAppend.set(true);
}
return null;
}).when(listener).onParagraphOutputUpdated(
"note", "second", 0, null, InterpreterResult.Type.TEXT, "replacement");
doAnswer(call -> {
if (appendCallbackActive.get() || !appendCompleted.get()) {
boundaryOverlappedAppend.set(true);
}
return null;
}).when(listener).onParagraphOutputClear("note", "second", null);
doAnswer(call -> {
if (appendCallbackActive.get() || !appendCompleted.get()) {
boundaryOverlappedAppend.set(true);
}
return null;
}).when(listener).checkpointOutput("note", "second");

ExecutorService callers = Executors.newSingleThreadExecutor();
try {
server.appendOutput(new OutputAppendEvent("note", "first", 0, "old", null, null));
dispatcherOf(server).flush();
assertTrue(entered.await(5, TimeUnit.SECONDS));
Future<?> boundary = callers.submit(() -> {
started.countDown();
callerThread.set(Thread.currentThread());
callBoundary(server, "note", "second", operation);
rpcReturnedEarly.set(!appendCompleted.get());
return null;
});
assertTrue(started.await(5, TimeUnit.SECONDS));
assertThrows(TimeoutException.class, () -> boundary.get(100, TimeUnit.MILLISECONDS));
await().atMost(5, TimeUnit.SECONDS).until(() -> boundary.isDone()
|| (callerThread.get() != null
&& callerThread.get().getState() == Thread.State.WAITING));

assertFalse(appendCompleted.get());
assertFalse(boundary.isDone());

release.countDown();
boundary.get(5, TimeUnit.SECONDS);

assertFalse(rpcReturnedEarly.get(), "RPC must wait for the earlier append");
assertFalse(boundaryOverlappedAppend.get(), "Boundary must wait for append completion");

InOrder order = inOrder(listener);
order.verify(listener).onParagraphOutputAppend("note", "first", 0, null, "old");
if ("UPDATE".equals(operation)) {
Expand Down Expand Up @@ -259,47 +300,73 @@ void laterAppendCannotOvertakeBoundaryCallback(String operation) throws Exceptio
RemoteInterpreterEventServer server = serverWithListener(listener);
CountDownLatch entered = new CountDownLatch(1);
CountDownLatch release = new CountDownLatch(1);
AtomicBoolean boundaryCallbackActive = new AtomicBoolean();
AtomicBoolean boundaryCompleted = new AtomicBoolean();
AtomicBoolean rpcReturnedEarly = new AtomicBoolean();
AtomicBoolean appendOverlappedBoundary = new AtomicBoolean();

doAnswer(call -> {
boundaryCallbackActive.set(true);
entered.countDown();
assertTrue(release.await(5, TimeUnit.SECONDS));
try {
assertTrue(release.await(5, TimeUnit.SECONDS));
} finally {
boundaryCallbackActive.set(false);
}
return null;
}).when(listener).onParagraphOutputClear("note", "para", null);
doAnswer(call -> {
boundaryCallbackActive.set(true);
entered.countDown();
assertTrue(release.await(5, TimeUnit.SECONDS));
boundaryCompleted.set(true);
try {
assertTrue(release.await(5, TimeUnit.SECONDS));
} finally {
boundaryCallbackActive.set(false);
boundaryCompleted.set(true);
}
return null;
}).when(listener).checkpointOutput("note", "para");
doAnswer(call -> {
boundaryCallbackActive.set(true);
entered.countDown();
assertTrue(release.await(5, TimeUnit.SECONDS));
boundaryCompleted.set(true);
try {
assertTrue(release.await(5, TimeUnit.SECONDS));
} finally {
boundaryCallbackActive.set(false);
boundaryCompleted.set(true);
}
return null;
}).when(listener).onParagraphOutputUpdated("note", "para", 0, null, InterpreterResult.Type.TEXT, "replacement");
doAnswer(call -> {
if (!boundaryCompleted.get()) {
if (boundaryCallbackActive.get() || !boundaryCompleted.get()) {
appendOverlappedBoundary.set(true);
}
return null;
}).when(listener).onParagraphOutputAppend("note", "para", 0, null, "later");

ExecutorService callers = Executors.newSingleThreadExecutor();
try {
Future<?> boundary = callers.submit(() -> {
callBoundary(server, "note", "para", operation);
rpcReturnedEarly.set(!boundaryCompleted.get());
return null;
});
assertTrue(entered.await(5, TimeUnit.SECONDS));
server.appendOutput(new OutputAppendEvent("note", "para", 0, "later", null, null));
dispatcherOf(server).flush();

verify(listener, never()).onParagraphOutputAppend("note", "para", 0, null, "later");
assertThrows(TimeoutException.class, () -> boundary.get(100, TimeUnit.MILLISECONDS));
assertFalse(boundaryCompleted.get());
assertFalse(boundary.isDone());

release.countDown();
boundary.get(5, TimeUnit.SECONDS);
server.checkpointOutput("note", "drained");

assertTrue(boundaryCompleted.get());
assertFalse(rpcReturnedEarly.get(), "RPC must wait for boundary completion");
assertFalse(appendOverlappedBoundary.get(), "Append must wait for the entire boundary");

InOrder order = inOrder(listener);
if ("UPDATE".equals(operation)) {
order.verify(listener).onParagraphOutputUpdated("note", "para", 0, null, InterpreterResult.Type.TEXT, "replacement");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.zeppelin.conf.ZeppelinConfiguration;
Expand Down Expand Up @@ -361,40 +362,53 @@ void blockedCheckpointsDoNotOccupyOutputWorkersOrReleaseTheirNotes() throws Exce
CountDownLatch release = new CountDownLatch(1);
CountDownLatch unrelatedOutput = new CountDownLatch(1);
CountDownLatch laterOutput = new CountDownLatch(1);
AtomicReference<Future<Void>> firstCheckpoint = new AtomicReference<>();
AtomicBoolean appendOverlappedCheckpoint = new AtomicBoolean();

doAnswer(call -> {
entered.countDown();
awaitIgnoringInterrupt(release);
return null;
}).when(listener).checkpointOutput(anyString(), anyString());

doAnswer(call -> {
unrelatedOutput.countDown();
return null;
}).when(listener).onParagraphOutputAppend("C", "para", 0, null, "unrelated");
doAnswer(call -> {
appendOverlappedCheckpoint.set(!firstCheckpoint.get().isDone());
laterOutput.countDown();
return null;
}).when(listener).onParagraphOutputAppend("A", "para", 0, null, "later");

try (ParagraphOutputDispatcher dispatcher = new ParagraphOutputDispatcher(listener, 2)) {
Future<Void> first = dispatcher.checkpointOutput("A", "para");
firstCheckpoint.set(first);
Future<Void> second = dispatcher.checkpointOutput("B", "para");
assertTrue(entered.await(5, TimeUnit.SECONDS));

dispatcher.appendOutput("A", "para", 0, null, "later");
dispatcher.appendOutput("C", "para", 0, null, "unrelated");
Future<Void> third = dispatcher.checkpointOutput("C", "para");
dispatcher.flush();

assertTrue(unrelatedOutput.await(5, TimeUnit.SECONDS));
assertFalse(laterOutput.await(100, TimeUnit.MILLISECONDS));
assertEquals(1, laterOutput.getCount());
assertFalse(first.isDone());
assertFalse(second.isDone());
assertFalse(third.isDone());
verify(listener, never()).checkpointOutput("C", "para");

release.countDown();
first.get(5, TimeUnit.SECONDS);
second.get(5, TimeUnit.SECONDS);
third.get(5, TimeUnit.SECONDS);
assertTrue(laterOutput.await(5, TimeUnit.SECONDS));

assertFalse(appendOverlappedCheckpoint.get(), "Append must wait for checkpoint release");
verify(listener).checkpointOutput("A", "para");
verify(listener).checkpointOutput("B", "para");

InOrder order = inOrder(listener);
order.verify(listener).checkpointOutput("A", "para");
order.verify(listener).onParagraphOutputAppend("A", "para", 0, null, "later");
Expand Down
Loading