Skip to content

Commit f164e31

Browse files
committed
[ZEPPELIN-6709] Clarify output ordering and strengthen boundary tests
1 parent 7a2e439 commit f164e31

4 files changed

Lines changed: 99 additions & 17 deletions

File tree

‎conf/zeppelin-site.xml.template‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -443,8 +443,9 @@
443443
<name>zeppelin.interpreter.output.worker.count</name>
444444
<value>4</value>
445445
<description>Number of workers in each server pool: one pool delivers paragraph output, and the
446-
other saves output checkpoints. Output for a note stays in order, while different notes can be
447-
processed concurrently. Checkpoint saves do not occupy output workers.</description>
446+
other saves output checkpoints. Each paragraph's output stays in order, and updates or checkpoints
447+
never overtake earlier appends. Different notes can be processed concurrently. Checkpoint saves
448+
do not occupy output workers.</description>
448449
</property>
449450

450451
<property>

‎docs/setup/operation/configuration.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -374,7 +374,7 @@ Sources descending by priority:
374374
<td><h6 class="properties">ZEPPELIN_INTERPRETER_OUTPUT_WORKER_COUNT</h6></td>
375375
<td><h6 class="properties">zeppelin.interpreter.output.worker.count</h6></td>
376376
<td>4</td>
377-
<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>
377+
<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>
378378
</tr>
379379
<tr>
380380
<td><h6 class="properties">ZEPPELIN_INTERPRETER_OUTPUT_EVENTS_PER_BATCH</h6></td>

‎zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerTest.java‎

Lines changed: 80 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
*/
1717
package org.apache.zeppelin.interpreter;
1818

19+
import static org.awaitility.Awaitility.await;
1920
import static org.junit.jupiter.api.Assertions.assertFalse;
2021
import static org.junit.jupiter.api.Assertions.assertThrows;
2122
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -45,7 +46,6 @@
4546
import java.util.concurrent.Executors;
4647
import java.util.concurrent.Future;
4748
import java.util.concurrent.TimeUnit;
48-
import java.util.concurrent.TimeoutException;
4949
import java.util.concurrent.atomic.AtomicBoolean;
5050
import java.util.concurrent.atomic.AtomicReference;
5151
import org.apache.zeppelin.conf.ZeppelinConfiguration;
@@ -174,26 +174,67 @@ void sameNoteBoundaryWaitsForInFlightAppend(String operation) throws Exception {
174174
RemoteInterpreterEventServer server = serverWithListener(listener);
175175
CountDownLatch entered = new CountDownLatch(1);
176176
CountDownLatch release = new CountDownLatch(1);
177-
CountDownLatch started = new CountDownLatch(1);
177+
AtomicReference<Thread> callerThread = new AtomicReference<>();
178+
AtomicBoolean appendCallbackActive = new AtomicBoolean();
179+
AtomicBoolean appendCompleted = new AtomicBoolean();
180+
AtomicBoolean boundaryOverlappedAppend = new AtomicBoolean();
181+
AtomicBoolean rpcReturnedEarly = new AtomicBoolean();
182+
178183
doAnswer(invocation -> {
184+
appendCallbackActive.set(true);
179185
entered.countDown();
180-
assertTrue(release.await(5, TimeUnit.SECONDS));
186+
try {
187+
assertTrue(release.await(5, TimeUnit.SECONDS));
188+
} finally {
189+
appendCallbackActive.set(false);
190+
appendCompleted.set(true);
191+
}
181192
return null;
182193
}).when(listener).onParagraphOutputAppend("note", "first", 0, null, "old");
194+
doAnswer(call -> {
195+
if (appendCallbackActive.get() || !appendCompleted.get()) {
196+
boundaryOverlappedAppend.set(true);
197+
}
198+
return null;
199+
}).when(listener).onParagraphOutputUpdated(
200+
"note", "second", 0, null, InterpreterResult.Type.TEXT, "replacement");
201+
doAnswer(call -> {
202+
if (appendCallbackActive.get() || !appendCompleted.get()) {
203+
boundaryOverlappedAppend.set(true);
204+
}
205+
return null;
206+
}).when(listener).onParagraphOutputClear("note", "second", null);
207+
doAnswer(call -> {
208+
if (appendCallbackActive.get() || !appendCompleted.get()) {
209+
boundaryOverlappedAppend.set(true);
210+
}
211+
return null;
212+
}).when(listener).checkpointOutput("note", "second");
213+
183214
ExecutorService callers = Executors.newSingleThreadExecutor();
184215
try {
185216
server.appendOutput(new OutputAppendEvent("note", "first", 0, "old", null, null));
186217
dispatcherOf(server).flush();
187218
assertTrue(entered.await(5, TimeUnit.SECONDS));
188219
Future<?> boundary = callers.submit(() -> {
189-
started.countDown();
220+
callerThread.set(Thread.currentThread());
190221
callBoundary(server, "note", "second", operation);
222+
rpcReturnedEarly.set(!appendCompleted.get());
191223
return null;
192224
});
193-
assertTrue(started.await(5, TimeUnit.SECONDS));
194-
assertThrows(TimeoutException.class, () -> boundary.get(100, TimeUnit.MILLISECONDS));
225+
await().atMost(5, TimeUnit.SECONDS).until(() -> boundary.isDone()
226+
|| (callerThread.get() != null
227+
&& callerThread.get().getState() == Thread.State.WAITING));
228+
229+
assertFalse(appendCompleted.get());
230+
assertFalse(boundary.isDone());
231+
195232
release.countDown();
196233
boundary.get(5, TimeUnit.SECONDS);
234+
235+
assertFalse(rpcReturnedEarly.get(), "RPC must wait for the earlier append");
236+
assertFalse(boundaryOverlappedAppend.get(), "Boundary must wait for append completion");
237+
197238
InOrder order = inOrder(listener);
198239
order.verify(listener).onParagraphOutputAppend("note", "first", 0, null, "old");
199240
if ("UPDATE".equals(operation)) {
@@ -259,47 +300,73 @@ void laterAppendCannotOvertakeBoundaryCallback(String operation) throws Exceptio
259300
RemoteInterpreterEventServer server = serverWithListener(listener);
260301
CountDownLatch entered = new CountDownLatch(1);
261302
CountDownLatch release = new CountDownLatch(1);
303+
AtomicBoolean boundaryCallbackActive = new AtomicBoolean();
262304
AtomicBoolean boundaryCompleted = new AtomicBoolean();
305+
AtomicBoolean rpcReturnedEarly = new AtomicBoolean();
263306
AtomicBoolean appendOverlappedBoundary = new AtomicBoolean();
307+
264308
doAnswer(call -> {
309+
boundaryCallbackActive.set(true);
265310
entered.countDown();
266-
assertTrue(release.await(5, TimeUnit.SECONDS));
311+
try {
312+
assertTrue(release.await(5, TimeUnit.SECONDS));
313+
} finally {
314+
boundaryCallbackActive.set(false);
315+
}
267316
return null;
268317
}).when(listener).onParagraphOutputClear("note", "para", null);
269318
doAnswer(call -> {
319+
boundaryCallbackActive.set(true);
270320
entered.countDown();
271-
assertTrue(release.await(5, TimeUnit.SECONDS));
272-
boundaryCompleted.set(true);
321+
try {
322+
assertTrue(release.await(5, TimeUnit.SECONDS));
323+
} finally {
324+
boundaryCallbackActive.set(false);
325+
boundaryCompleted.set(true);
326+
}
273327
return null;
274328
}).when(listener).checkpointOutput("note", "para");
275329
doAnswer(call -> {
330+
boundaryCallbackActive.set(true);
276331
entered.countDown();
277-
assertTrue(release.await(5, TimeUnit.SECONDS));
278-
boundaryCompleted.set(true);
332+
try {
333+
assertTrue(release.await(5, TimeUnit.SECONDS));
334+
} finally {
335+
boundaryCallbackActive.set(false);
336+
boundaryCompleted.set(true);
337+
}
279338
return null;
280339
}).when(listener).onParagraphOutputUpdated("note", "para", 0, null, InterpreterResult.Type.TEXT, "replacement");
281340
doAnswer(call -> {
282-
if (!boundaryCompleted.get()) {
341+
if (boundaryCallbackActive.get() || !boundaryCompleted.get()) {
283342
appendOverlappedBoundary.set(true);
284343
}
285344
return null;
286345
}).when(listener).onParagraphOutputAppend("note", "para", 0, null, "later");
346+
287347
ExecutorService callers = Executors.newSingleThreadExecutor();
288348
try {
289349
Future<?> boundary = callers.submit(() -> {
290350
callBoundary(server, "note", "para", operation);
351+
rpcReturnedEarly.set(!boundaryCompleted.get());
291352
return null;
292353
});
293354
assertTrue(entered.await(5, TimeUnit.SECONDS));
294355
server.appendOutput(new OutputAppendEvent("note", "para", 0, "later", null, null));
295356
dispatcherOf(server).flush();
357+
296358
verify(listener, never()).onParagraphOutputAppend("note", "para", 0, null, "later");
297-
assertThrows(TimeoutException.class, () -> boundary.get(100, TimeUnit.MILLISECONDS));
359+
assertFalse(boundaryCompleted.get());
360+
assertFalse(boundary.isDone());
361+
298362
release.countDown();
299363
boundary.get(5, TimeUnit.SECONDS);
300364
server.checkpointOutput("note", "drained");
365+
301366
assertTrue(boundaryCompleted.get());
367+
assertFalse(rpcReturnedEarly.get(), "RPC must wait for boundary completion");
302368
assertFalse(appendOverlappedBoundary.get(), "Append must wait for the entire boundary");
369+
303370
InOrder order = inOrder(listener);
304371
if ("UPDATE".equals(operation)) {
305372
order.verify(listener).onParagraphOutputUpdated("note", "para", 0, null, InterpreterResult.Type.TEXT, "replacement");

‎zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/ParagraphOutputDispatcherTest.java‎

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,6 +53,7 @@
5353
import java.util.concurrent.Future;
5454
import java.util.concurrent.ScheduledExecutorService;
5555
import java.util.concurrent.TimeUnit;
56+
import java.util.concurrent.atomic.AtomicBoolean;
5657
import java.util.concurrent.atomic.AtomicInteger;
5758
import java.util.concurrent.atomic.AtomicReference;
5859
import org.apache.zeppelin.conf.ZeppelinConfiguration;
@@ -361,40 +362,53 @@ void blockedCheckpointsDoNotOccupyOutputWorkersOrReleaseTheirNotes() throws Exce
361362
CountDownLatch release = new CountDownLatch(1);
362363
CountDownLatch unrelatedOutput = new CountDownLatch(1);
363364
CountDownLatch laterOutput = new CountDownLatch(1);
365+
AtomicReference<Future<Void>> firstCheckpoint = new AtomicReference<>();
366+
AtomicBoolean appendOverlappedCheckpoint = new AtomicBoolean();
367+
364368
doAnswer(call -> {
365369
entered.countDown();
366370
awaitIgnoringInterrupt(release);
367371
return null;
368372
}).when(listener).checkpointOutput(anyString(), anyString());
373+
369374
doAnswer(call -> {
370375
unrelatedOutput.countDown();
371376
return null;
372377
}).when(listener).onParagraphOutputAppend("C", "para", 0, null, "unrelated");
373378
doAnswer(call -> {
379+
appendOverlappedCheckpoint.set(!firstCheckpoint.get().isDone());
374380
laterOutput.countDown();
375381
return null;
376382
}).when(listener).onParagraphOutputAppend("A", "para", 0, null, "later");
383+
377384
try (ParagraphOutputDispatcher dispatcher = new ParagraphOutputDispatcher(listener, 2)) {
378385
Future<Void> first = dispatcher.checkpointOutput("A", "para");
386+
firstCheckpoint.set(first);
379387
Future<Void> second = dispatcher.checkpointOutput("B", "para");
380388
assertTrue(entered.await(5, TimeUnit.SECONDS));
389+
381390
dispatcher.appendOutput("A", "para", 0, null, "later");
382391
dispatcher.appendOutput("C", "para", 0, null, "unrelated");
383392
Future<Void> third = dispatcher.checkpointOutput("C", "para");
384393
dispatcher.flush();
394+
385395
assertTrue(unrelatedOutput.await(5, TimeUnit.SECONDS));
386-
assertFalse(laterOutput.await(100, TimeUnit.MILLISECONDS));
396+
assertEquals(1, laterOutput.getCount());
387397
assertFalse(first.isDone());
388398
assertFalse(second.isDone());
389399
assertFalse(third.isDone());
390400
verify(listener, never()).checkpointOutput("C", "para");
401+
391402
release.countDown();
392403
first.get(5, TimeUnit.SECONDS);
393404
second.get(5, TimeUnit.SECONDS);
394405
third.get(5, TimeUnit.SECONDS);
395406
assertTrue(laterOutput.await(5, TimeUnit.SECONDS));
407+
408+
assertFalse(appendOverlappedCheckpoint.get(), "Append must wait for checkpoint release");
396409
verify(listener).checkpointOutput("A", "para");
397410
verify(listener).checkpointOutput("B", "para");
411+
398412
InOrder order = inOrder(listener);
399413
order.verify(listener).checkpointOutput("A", "para");
400414
order.verify(listener).onParagraphOutputAppend("A", "para", 0, null, "later");

0 commit comments

Comments
 (0)