Skip to content

[ZEPPELIN-6340] Add ZeppelinEventBus and update NotebookServer to handle NoteRemoveEvent - #5441

Open
seung-00 wants to merge 3 commits into
apache:masterfrom
seung-00:ZEPPELIN-6340
Open

seung-00 wants to merge 3 commits into
apache:masterfrom
seung-00:ZEPPELIN-6340

Conversation

@seung-00

@seung-00 seung-00 commented Aug 29, 2026 •

Copy link
Copy Markdown
Contributor

What is this PR for?

This PR is the first step toward migrating Zeppelin’s listener-based event handling to an EventBus model, as outlined in the proposal. It sets up the EventBus infrastructure and migrates NoteRemoveEvent handling to it.

Key changes:

  • Added the ZeppelinEventBus class
  • Added RxJava 3 as a dependency of the zeppelin-zengine module
  • Added the zeppelin.eventbus.enabled feature flag
  • Updated NotebookServer to handle NoteRemoveEvent through ZeppelinEventBus when the feature flag is enabled

This PR supersedes #5085. The original branch had fallen behind the current codebase, and the implementation has since changed, so I opened a new PR.

What type of PR is it?

Improvement

Todos

  • Refactor other NoteEventListener methods and additional listeners to EventBus

What is the Jira issue?

How should this be tested?

  • Run the existing test suite against the updated code
  • Verify the note deletion flow through the UI

Screenshots (if appropriate)

Questions:

  • Does the license files need to update? N
  • Is there breaking changes for older versions? N
  • Does this needs documentation? N

@seung-00
seung-00 force-pushed the ZEPPELIN-6340 branch 4 times, most recently from 73168af to 4ff5711 Compare August 29, 2026 14:46
@seung-00
seung-00 force-pushed the ZEPPELIN-6340 branch 3 times, most recently from 58118cb to 058654a Compare September 7, 2026 06:10

@tbonelee tbonelee left a comment

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.

A few thoughts on subscription registration, since 6341, 6346 and 6348 will likely follow this PR as their model.

observe(NoteEvent.class) receives every subtype of NoteEvent. Once 6348 adds NoteCreatedEvent, this subscription gets it too and falls through to the else branch in handleNoteEvent, logging "Unknown event type" for what is normal behavior. Each new event also means another ticket editing the same if-else chain.

The try/catch in the lambda is load-bearing in a way that is easy to miss. LambdaObserver disposes the subscription when the consumer throws and forwards the error to RxJavaPlugins.onError, which by default only reaches stderr, so a dead subscription leaves no trace in the logs.

Would something like subscribe(Class<T>, Consumer<T>) on EventBus work here? Subscribing by concrete type lets ofType() do the filtering instead of instanceof, and the try/catch sits in the bus once however many subscriptions there are. One caveat: Consumer would need to be io.reactivex.rxjava3.functions.Consumer, since java.util.function.Consumer cannot throw checked exceptions and handlers would go on catching IOException themselves.

@seung-00

seung-00 commented Oct 5, 2026

Copy link
Copy Markdown
Contributor Author

@tbonelee
Thanks for the suggestion. I've applied it.

@seung-00
seung-00 requested a review from tbonelee October 5, 2026 16:14

@tbonelee tbonelee left a comment •

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.

Three things, labelled so you can skip what isn't worth your time.

issue (blocking): catch (Exception) leaves Error uncovered.

An Error from a handler escapes that catch and reaches LambdaObserver, which calls Exceptions.throwIfFatal first. VirtualMachineError, LinkageError and ThreadDeath get rethrown, everything else reaches dispose() and RxJavaPlugins.onError. So the app keeps running, that subscription is gone, and only stderr shows it. I ran the subscribe path as it stands against 3.1.12 to check.

catch (Throwable t) with Exceptions.throwIfFatal(t) on the first line would cover it, same shape LambdaObserver uses. Probably worth doing before this lands, since 6341, 6346 and 6348 will copy whatever shape ends up here.

AssertionError is in that second group, so a test that asserts inside the handler lambda passes even when the assertion fails. Maybe a third case in testContinueAfterException throwing an Error directly?

suggestion (docs): a handler can run on a thread other than the one that called post.

toSerialized() lets a publisher already inside the bus drain events other threads queued, so with two overlapping removals thread A executes B's event while Shiro's ThreadContext still holds A. Probably fine today, since handleNoteRemove works off the Note it is handed, but reaching for SecurityUtils.getSubject() in a later handler would be easy.

Might be worth a line of javadoc on subscribe saying to take what you need from the event, which I assume is why NoteEvent carries getSubject().

note (for the follow-up tickets): with observe() gone there is no place to put an operator between the subject and the handler.

The case that worries me is LuceneSearch / EmbeddingSearch. Both extend NoteEventAsyncListener, which pushes every event onto its own single-thread executor. Migrating them to subscribe(Class, Consumer) as-is puts Lucene indexing and embedding inference on the calling thread. Nothing throws, nothing logs, the async handoff is just gone.

Maybe a subscribe(Class, Scheduler, Consumer) overload, though I haven't thought hard enough about what the right extension would be.

@seung-00

seung-00 commented Oct 11, 2026 •

Copy link
Copy Markdown
Contributor Author

@tbonelee

issue (blocking): catch (Exception) leaves Error uncovered.

In Java, the general recommendation is not to catch Error. But RxJava reports some errors through onError rather than throwing them, because throwing an error cannot reliably notify the caller in asynchronous execution.

I’d prefer to preserve RxJava’s intended behavior for now. Once we have concrete use cases requiring keep subscription, we can explicitly handle those errors and continue processing. I think Errors classified as nonfatal by RxJava can still be unsafe to swallow.

suggestion (docs): a handler can run on a thread other than the one that called post.

Agreed. I’ve added docs.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants