Repository navigation
Conversation
73168af to
4ff5711
Compare
58118cb to
058654a
Compare
tbonelee
left a comment
There was a problem hiding this comment.
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.
058654a to
8afc86b
Compare
|
@tbonelee |
There was a problem hiding this comment.
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.
In Java, the general recommendation is not to catch
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.
Agreed. I’ve added docs. |
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
NoteRemoveEventhandling to it.Key changes:
ZeppelinEventBusclasszeppelin-zenginemodulezeppelin.eventbus.enabledfeature flagNotebookServerto handleNoteRemoveEventthroughZeppelinEventBuswhen the feature flag is enabledThis 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
What is the Jira issue?
How should this be tested?
Screenshots (if appropriate)
Questions: