This is an automated email from the ASF dual-hosted git repository.
tbonelee pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/zeppelin.git
The following commit(s) were added to refs/heads/master by this push:
new b5d927290b [ZEPPELIN-6700] Convert a broadcast message to JSON once
instead of once per connection
b5d927290b is described below
commit b5d927290b656e5939b08fbf2684aa0afa5bb609
Author: dae won <[email protected]>
AuthorDate: Thu Sep 10 00:56:39 2026 +0900
[ZEPPELIN-6700] Convert a broadcast message to JSON once instead of once
per connection
### What is this PR for?
`ConnectionManager` converts a `Message` to JSON using `serializeMessage`
and writes the result to a socket. A broadcast sends one message to many
connections, so the two steps have different multiplicity: the JSON is one
value, the writes are many. Five call sites currently perform the conversion
inside the per-connection loop, causing it to run once per connection instead
of once per message:
* `broadcast(Message)`, inside the `synchronized (connectedSockets)` block
* `broadcast(String, Message)`
* `broadcastToWatchers(String, String, Message)`, itself reached once per
broadcast from `broadcast`, `broadcastExcept`, and `unicast`
* `broadcastExcept(String, Message, NotebookSocket)`, the path used by
collaborative patches, Angular object updates, and spell results
* `unicast(Message, NotebookSocket)`, reached once per connection from
`multicastToUser`, which serves both note-list updates and personalized-mode
paragraphs
The message is not mutated inside these loops, and `gson.toJson` is
deterministic. Every iteration therefore produces a byte-identical string, and
all but one are discarded.
There is no functional bug being fixed here. The bytes on the wire and
their order are already correct, and on a note with a single connection the
current form costs nothing extra. What it does is repeat work in proportion to
the number of connections attached to a note, and three paths make that
repetition routine rather than occasional:
1. A paragraph state change broadcasts the whole `Paragraph` including its
output, which `zeppelin.interpreter.output.limit` caps at 100 KB by default.
This is the path where a single redundant conversion is expensive.
2. While a paragraph runs, streaming output is broadcast once per
`AppendOutputRunner` flush. Each message carries only the appended chunk and
the runner coalesces writes over a 100 ms window, so these are small but
numerous.
3. In collaborative mode, every edit to a paragraph broadcasts a patch to
the other connections on the note.
The features that make a note worth sharing are the ones that pay for this
most.
This PR converts once per broadcast at those five sites and writes that one
string to every connection. The resulting JSON and write order remain
unchanged. Two supporting changes come with it:
* `broadcastToWatchers` returns early when no watcher is attached, so
hoisting the conversion out of its loop does not introduce work in the common
case where the loop body never ran.
* `multicastToUser` now sends directly rather than delegating to `unicast`,
because `unicast` bundles the conversion with a watcher broadcast. The watcher
broadcast stays inside the loop, so watchers receive the same messages they do
today.
One behavioral detail is worth flagging for review. Hoisting moves the
conversion out of the per-connection `try` that exists to catch `IOException`
from `NotebookSocket.send`. A conversion failure was logged and skipped per
connection before, and now propagates to the caller. That handler was
incidental to `send` rather than an intentional contract for the conversion,
and it never caught `StackOverflowError`, which is the likely failure mode for
a cyclic object graph.
`NotebookServer` is untouched. Its own `serializeMessage` is used for
single sends, and the one loop there that converts per iteration builds a
different message each time, so there is nothing to hoist.
### What type of PR is it?
Improvement
### Todos
* [x] - Convert once per broadcast in `broadcast(Message)`,
`broadcast(String, Message)`, `broadcastExcept`, and `broadcastToWatchers`
* [x] - Return early from `broadcastToWatchers` when no watcher is attached
* [x] - Send directly from `multicastToUser` instead of going through
`unicast`, keeping the watcher broadcast inside the loop
* [x] - Add tests asserting that a single broadcast to multiple connections
converts once and delivers identical payloads
### What is the Jira issue?
* [ZEPPELIN-6700](https://issues.apache.org/jira/browse/ZEPPELIN-6700)
### How should this be tested?
Three unit tests were added to `ConnectionManagerTest`. They subclass
`ConnectionManager` to count how often a broadcast converts its message and
assert that every connection receives an identical payload.
```bash
./mvnw package -pl zeppelin-server --am \
-Dtest=ConnectionManagerTest,NotebookServerTest -DfailIfNoTests=false
```
### Screenshots (if appropriate)
N/A
### Questions:
* Does the license files need to update? No
* Is there breaking changes for older versions? No
* Does this needs documentation? No
Closes #5462 from big-cir/ZEPPELIN-6700.
Signed-off-by: ChanHo Lee <[email protected]>
---
.../apache/zeppelin/socket/ConnectionManager.java | 32 ++++---
.../zeppelin/socket/ConnectionManagerTest.java | 98 ++++++++++++++++++++++
2 files changed, 120 insertions(+), 10 deletions(-)
diff --git
a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/ConnectionManager.java
b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/ConnectionManager.java
index 6b13613cce..4c2b4d463c 100644
---
a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/ConnectionManager.java
+++
b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/ConnectionManager.java
@@ -229,10 +229,11 @@ public class ConnectionManager {
}
public void broadcast(Message m) {
+ String serializedMsg = serializeMessage(m);
synchronized (connectedSockets) {
for (NotebookSocket ns : connectedSockets) {
try {
- ns.send(serializeMessage(m));
+ ns.send(serializedMsg);
} catch (IOException | RuntimeException e) {
LOGGER.error("Send error: {}", m, e);
}
@@ -251,9 +252,10 @@ public class ConnectionManager {
socketsToBroadcast = new ArrayList<>(sockets);
}
LOGGER.debug("SEND >> {}", m);
+ String serializedMsg = serializeMessage(m);
for (NotebookSocket conn : socketsToBroadcast) {
try {
- conn.send(serializeMessage(m));
+ conn.send(serializedMsg);
} catch (IOException | RuntimeException e) {
LOGGER.error("socket error", e);
}
@@ -262,14 +264,17 @@ public class ConnectionManager {
private void broadcastToWatchers(String noteId, String subject, Message
message) {
synchronized (watcherSockets) {
+ if (watcherSockets.isEmpty()) {
+ return;
+ }
+ String watcherMsg = WatcherMessage.builder(noteId)
+ .subject(subject)
+ .message(serializeMessage(message))
+ .build()
+ .toJson();
for (NotebookSocket watcher : watcherSockets) {
try {
- watcher.send(
- WatcherMessage.builder(noteId)
- .subject(subject)
- .message(serializeMessage(message))
- .build()
- .toJson());
+ watcher.send(watcherMsg);
} catch (IOException | RuntimeException e) {
LOGGER.error("Cannot broadcast message to watcher", e);
}
@@ -289,12 +294,13 @@ public class ConnectionManager {
}
LOGGER.debug("SEND >> {}", m);
+ String serializedMsg = serializeMessage(m);
for (NotebookSocket conn : socketsToBroadcast) {
if (exclude.equals(conn)) {
continue;
}
try {
- conn.send(serializeMessage(m));
+ conn.send(serializedMsg);
} catch (IOException | RuntimeException e) {
LOGGER.error("socket error", e);
}
@@ -340,8 +346,14 @@ public class ConnectionManager {
return;
}
+ String serializedMsg = serializeMessage(m);
for (NotebookSocket conn : connections) {
- unicast(m, conn);
+ try {
+ conn.send(serializedMsg);
+ } catch (IOException | RuntimeException e) {
+ LOGGER.error("socket error", e);
+ }
+ broadcastToWatchers(StringUtils.EMPTY, StringUtils.EMPTY, m);
}
}
diff --git
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/ConnectionManagerTest.java
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/ConnectionManagerTest.java
index 562d065894..8d0ece1451 100644
---
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/ConnectionManagerTest.java
+++
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/ConnectionManagerTest.java
@@ -22,9 +22,13 @@ import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
+import java.io.IOException;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.List;
@@ -34,12 +38,15 @@ import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
+import org.apache.zeppelin.common.Message;
import org.apache.zeppelin.conf.ZeppelinConfiguration;
import org.apache.zeppelin.notebook.AuthorizationService;
import org.apache.zeppelin.util.WatcherSecurityKey;
import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
class ConnectionManagerTest {
@@ -310,4 +317,95 @@ class ConnectionManagerTest {
assertDoesNotThrow(() -> manager.removeUserConnection("", socket));
assertTrue(manager.userSocketMap.isEmpty());
}
+
+ @Test
+ void broadcastToNoteSerializesMessageOnce() throws IOException {
+ CountingConnectionManager manager = newCountingManager();
+ NotebookSocket first = mock(NotebookSocket.class);
+ NotebookSocket second = mock(NotebookSocket.class);
+ NotebookSocket third = mock(NotebookSocket.class);
+ manager.addNoteConnection("note1", first);
+ manager.addNoteConnection("note1", second);
+ manager.addNoteConnection("note1", third);
+ // adding a connection broadcasts the collaborative mode status, which is
not under test
+ manager.serializeCount.set(0);
+ clearInvocations(first, second, third);
+
+ manager.broadcast("note1", new Message(Message.OP.NOTE).put("note",
"payload"));
+
+ assertEquals(1, manager.serializeCount.get());
+ assertEquals(payloadSentTo(first), payloadSentTo(second));
+ assertEquals(payloadSentTo(first), payloadSentTo(third));
+ }
+
+ @Test
+ void broadcastExceptSerializesMessageOnce() throws IOException {
+ CountingConnectionManager manager = newCountingManager();
+ NotebookSocket sender = mock(NotebookSocket.class);
+ NotebookSocket first = mock(NotebookSocket.class);
+ NotebookSocket second = mock(NotebookSocket.class);
+ manager.addNoteConnection("note1", sender);
+ manager.addNoteConnection("note1", first);
+ manager.addNoteConnection("note1", second);
+ manager.serializeCount.set(0);
+ clearInvocations(sender, first, second);
+
+ manager.broadcastExcept("note1",
+ new Message(Message.OP.PATCH_PARAGRAPH).put("patch", "payload"),
sender);
+
+ assertEquals(1, manager.serializeCount.get());
+ assertEquals(payloadSentTo(first), payloadSentTo(second));
+ verify(sender, times(0)).send(anyString());
+ }
+
+ @Test
+ void multicastToUserSerializesMessageOnceAndKeepsWatcherBroadcasts() throws
IOException {
+ CountingConnectionManager manager = newCountingManager();
+ NotebookSocket first = mock(NotebookSocket.class);
+ NotebookSocket second = mock(NotebookSocket.class);
+ manager.addUserConnection("testUser", first);
+ manager.addUserConnection("testUser", second);
+ NotebookSocket watcher = mock(NotebookSocket.class);
+
when(watcher.getHeader(WatcherSecurityKey.HTTP_HEADER)).thenReturn(WatcherSecurityKey.getKey());
+ manager.switchConnectionToWatcher(watcher);
+ manager.serializeCount.set(0);
+ clearInvocations(first, second, watcher);
+
+ manager.multicastToUser("testUser", new
Message(Message.OP.NOTES_INFO).put("notes", "payload"));
+
+ assertEquals(payloadSentTo(first), payloadSentTo(second));
+ // one conversion for the two connections, plus one per watcher broadcast
as before
+ assertEquals(3, manager.serializeCount.get());
+ verify(watcher, times(2)).send(anyString());
+ }
+
+ private static CountingConnectionManager newCountingManager() {
+ return new CountingConnectionManager(mock(AuthorizationService.class),
+ ZeppelinConfiguration.load());
+ }
+
+ private static String payloadSentTo(NotebookSocket socket) throws
IOException {
+ ArgumentCaptor<String> payload = ArgumentCaptor.forClass(String.class);
+ verify(socket).send(payload.capture());
+ return payload.getValue();
+ }
+
+ /**
+ * Counts how often a broadcast converts its message, which is what
distinguishes a single
+ * conversion reused across connections from one conversion per connection.
+ */
+ private static class CountingConnectionManager extends ConnectionManager {
+ private final AtomicInteger serializeCount = new AtomicInteger();
+
+ CountingConnectionManager(AuthorizationService authorizationService,
+ ZeppelinConfiguration zConf) {
+ super(authorizationService, zConf);
+ }
+
+ @Override
+ protected String serializeMessage(Message m) {
+ serializeCount.incrementAndGet();
+ return super.serializeMessage(m);
+ }
+ }
}