Repository navigation
[feat][client] PIP-445: Add Builder Methods to Create Message-based TableView - #24809
Conversation
…raw Pulsar message
…agement by adding retain() before storing message to rawMessages - Prevent premature message release that causes getRawMessage to return null keys. - Ensure message reference count is balanced by retaining message before storing and releasing old message after replacement. - Keep message release in finally block for safety.
There was a problem hiding this comment.
Thanks for the contribution!
Since this changes the Pulsar API, we'd need to first go through the PIP process to modify the public API.
Instead of adding a separate ConcurrentMap for storing the messages, I'd suggest going with an approach where the TableViewBuilder would have separate methods for creating a view for messages.
TableView<Message<T>> createForMessages() throws PulsarClientException;
CompletableFuture<TableView<Message<T>>> createForMessagesAsync() throws PulsarClientException;I'm not exactly sure if there's any obstacles with this approach, but it seems that it could be a better way forward so that the existing runtime behavior of TableView implementation wouldn't change (like it does with the current PR changes).
|
Hi @lhotari, Thanks for the review and the great suggestion about using the I completely agree that we should avoid impacting the runtime behavior for existing users. As you recommended, I've created a PIP for this API change, which you can find here: #24842 Let's continue the design discussion over on the PIP PR. |
@namest504 you seemed to ignore the suggestion in the comment. I'd suggest revisiting the PIP accordingly and renaming it. You can reply to the comment on the PIP PR, #24842 (review). |
Signed-off-by: namest504 <[email protected]>
Signed-off-by: namest504 <[email protected]>
Signed-off-by: namest504 <[email protected]>
… method Signed-off-by: namest504 <[email protected]>
Resolve conflicts by rebuilding the PIP-445 changes on top of current master: - Keep the original TableViewImpl class name; move the shared logic to AbstractTableViewImpl - Rename MessageMapperTableView to MessageMapperTableViewImpl - Drop the unused MessageTableView class and the import/FQN churn
…r argument - The mapper may retain the Message instance (e.g. Function.identity()), so the reader must not use pooled messages for mapped views; the payload buffer would otherwise be released or leaked. The classic TableView keeps using pooled messages and releasing them as before. - Reject a null mapper with IllegalArgumentException (failed future in the async variant) instead of failing later in the async pipeline. - Add unit tests for createMapped/createMappedAsync in TableViewBuilderImplTest.
|
@lhotari I've addressed the review comments and resolved the conflicts with master (rebuilt the refactoring on top of the current Gradle/slog codebase). Main changes since your last review:
The vote result for PIP-445 has been posted on the dev mailing list (3 binding +1s), and I've also updated the PIP document in #24842 to match the final class names. PTAL when you have a chance. |
…bleViewImpl FieldUtils.readDeclaredField only looks at the exact class, so reading the reader/compactionStrategy fields from TableViewImpl fails now that they live in AbstractTableViewImpl. Use readField, which traverses the class hierarchy.
Conflict in TableViewImpl.java: apache#26566 moved the lastReadPositions update in handleMessage to after the table is updated, while this PR moved handleMessage into AbstractTableViewImpl. Kept the PR side of TableViewImpl.java and ported the apache#26566 change into AbstractTableViewImpl.handleMessage.
lhotari
left a comment
There was a problem hiding this comment.
Sorry for the long silence — you asked for another look on 10 July and waited ten weeks for it. Congratulations on getting PIP-445 through the vote — the PIP document (#24842) is on master now too.
All four of my earlier points still check out: the naming is back to TableViewImpl / AbstractTableViewImpl / MessageMapperTableViewImpl, the import churn is gone, and MessageTableView is removed. Since then your fork merged upstream master again, and that merge correctly ports #26566 ("Complete table view refresh after applying messages") into the moved handleMessage in AbstractTableViewImpl: lastReadPositions is now recorded in each completion branch, after the value has been produced, instead of unconditionally before processing starts. That also closes a point I was going to raise — a throwing mapper leaving the read position advanced, so refreshAsync could report stale data as fresh — so I have dropped it.
I paid particular attention to the pooling decision, because it is the part most likely to go wrong, and it is right. TableViewImpl passes super(..., true) and MessageMapperTableViewImpl passes false; AbstractTableViewImpl releases in a finally only if (poolMessages), and the same flag drives .poolMessages(...) on the reader builder, so allocation and release cannot drift apart. I also diffed the classic path against master: it already used .poolMessages(true) with an unconditional release(), so classic TableView behaviour is preserved exactly. Your reasoning in the PR description — that the mapper may retain the Message, so the pooled buffer must not be reused — is the correct call.
What I would like changed is not the design but the contract around the mapper, and it matters more here than it would elsewhere because createMapped is new public API: once released, the behaviour is fixed. Four inline comments below; the first is the real question — what happens when a user's mapper throws? — and the other three are smaller.
That behaviour is pre-existing and I want to be fair about it. On master, TableViewImpl already routed lazy Message.getValue() schema-decode failures through the same handler. That is not something you introduced. What changes is the blast radius: a schema-decode failure is rare and roughly deterministic, whereas a mapper is arbitrary user code that this PR invites people to write. That is why I think the new API should either guard or document it rather than inherit it silently.
One smaller thing I considered and decided not to raise as a finding, because I could not substantiate the impact: the shared logger moved from Logger.get(TableViewImpl.class) to Logger.get(AbstractTableViewImpl.class), so classic-TableView log lines change category. Pulsar documents no compatibility guarantee for client logger names and I found no in-tree config keyed on the old one, so this is only worth doing if you think it is cheap — Logger.get(getClass()) in the constructor would keep per-subclass categories.
Nothing here questions PIP-445 or the shape of the change. It is close.
|
Thanks for the detailed summary, and sorry for the ten-week wait — that was on me, not on anything missing from your side. I have gone through it again at the current head. All four of my earlier points check out, and the pooling decision you describe is right: I diffed the classic path against master and it keeps My remaining comments are about the contract around the mapper rather than the design — mainly what should happen when a user's mapper throws, since |
…compaction strategies for mapped views - createMapped/createMappedAsync take a TableViewMessageMapper<T, V> instead of a Function: map() may throw, and onMappingError() reports failures (default: logged at ERROR by the view) - A message whose value extraction throws is skipped: the key keeps its previous value, listeners are not notified and the refresh position still advances - start() closes the reader when the initial replay fails instead of leaking it - Reject topicCompactionStrategyClassName for mapped views in both builder methods - Rename MessageMapperTableViewImpl to package-private MappedTableViewImpl - Document tombstone handling of empty payloads on the public API - Update PIP-445 to the revised API and semantics
# Conflicts: # pulsar-client/src/main/java/org/apache/pulsar/client/impl/TableViewImpl.java
|
@namest504 thank you for your persistence and patience in getting PIP-445 through the design discussion and the vote, and now finally into Pulsar — this took far longer than it should have, and that was largely on the review side. To get it over the line I pushed a final refactoring pass to your branch (maintainer edits are enabled) that revises the design slightly to tackle the error model: |
|
Nice work — the refactor is a faithful extraction. I diffed I ran Everything below is non-blocking; the PR already has an approval. Two points I'd like your read on, then some nits. 1. The
|
…iew constructor - readAllExistingMessages: a failed hasMessageAvailableAsync() now fails the replay instead of leaving create() waiting forever with the reader open - MappedTableViewImpl rejects a configured topicCompactionStrategyClassName itself, before the reader is created, and owns the message the builder reuses - Neutral log message for a skipped message, which may be a decode or a mapping failure - Document what an identity mapper retains and name the broker's create() consumer in the PIP-445 compatibility section
|
@david-streamlio thanks for the thorough second pass — the line-by-line diff against the base 1. 2. Nits
|
|
@lhotari Thank you for taking this over the line, and for the notes on each thread. I had been going back and forth on how a throwing mapper should be handled, and the direction you took makes the trade-off clear to me: skipping the message keeps a long-running view alive, and |
…ableView (apache#24809) Signed-off-by: namest504 <[email protected]> Co-authored-by: Lari Hotari <[email protected]>
Fixes: #24744
Main Issue: #24744
PIP: 445 #24842
Motivation
The current
TableViewAPI only exposes the deserialized message value, which limits access to essential message metadata like properties, event time, or the raw message object itself.This PR introduces a flexible, non-breaking mechanism to create value-mapped views over a topic. It replaces the initial proposal of adding a
getRawMessage()method, which would have performance implications for all users. Instead, it addscreateMappedmethods to theTableViewBuilder, allowing users to transform a fullMessage<T>into any custom objectVthat suits their needs. This provides maximum flexibility, from accessing the full raw message to creating custom, memory-efficient objects.Since arbitrary user code now runs inside the table view's message handling, the PR also defines what happens when that code fails. Previously an exception thrown while extracting a message's value escaped into the reader loop and was treated as a reader failure: during the initial replay it failed table view creation and leaked the reader, and while tailing the message was silently dropped with a log message that blamed the reader.
Modifications
TableViewMessageMapper<T, V>to the public client API: a functional interface withV map(Message<T>), which may throw, and adefault boolean onMappingError(Message<T>, Throwable)callback. The default callback returnsfalse, which makes the table view log the failure atERRORlevel with the topic, key and message id; an implementation can override it to report the failure its own way (metrics, signalling another thread, ...) and returntrueto suppress the log. Further callbacks can be added later asdefaultmethods without breaking implementations.createMappedandcreateMappedAsynctoTableViewBuilder: they take aTableViewMessageMapper<T, V>. Lambdas work directly, e.g.createMapped(msg -> msg.getProperty("region")), andcreateMapped(msg -> msg)builds aTableView<Message<T>>exposing the complete messages. Anullmapper is rejected withIllegalArgumentException(failed future in the async variant).nullreturn value from the mapper is also a tombstone; if the mapper throws, the message is skipped — the key keeps its previous value, listeners are not notified, a pendingrefresh()still completes — andonMappingErroris called. Exceptions thrown by the callback are logged and ignored.topicCompactionStrategyClassNamefor mapped views in both builder methods and, defensively, in theMappedTableViewImplconstructor before the reader is created: aTopicCompactionStrategycompares values of the topic's schema type, which a mapped view does not store, so the combination would fail with aClassCastExceptionon the reader thread.AbstractTableViewImpl<T, V>:TableViewImplkeeps its name and becomes a thin subclass, so the behavior for existing users (including message pooling and release) is unchanged, with one exception: a message whose lazily decoded value throws is now skipped and logged consistently in both the initial replay and the tailing phase, instead of failingcreate()in the former and being dropped silently in the latter. This also applies to the broker's owncreate()user, the extensible load balancer'sServiceUnitStateTableViewImpl: an undecodableServiceUnitStateDatarecord no longer prevents the channel from starting; it is skipped and logged atERROR, as it already was while tailing. The trade-off is documented in the PIP-445 compatibility section.MappedTableViewImpl<T, V>: implements the mapped views and delegates mapping failures to the mapper'sonMappingError. Message pooling is disabled for mapped views since the mapper may retain theMessageinstance; thecreateMappedjavadoc notes that a retainedMessagealso keeps its metadata, schema and connection reference, so a value object is preferable tomsg -> msgfor topics with many keys.start()fails: the reader created for the table view is no longer leaked when the initial replay fails. A failedhasMessageAvailableAsync()during the replay now failsstart()too, instead of leavingcreate()waiting forever with the reader open.Verifying this change
TableViewBuilderImplTest:testCreateMapped*(mapped view creation,nullmapper validation, and rejection of a configured compaction strategy for both the sync and async variants).MappedTableViewImplTestwith a mocked reader: a mapping failure skips the message and keeps the previous value, notifiesonMappingError, does not notify listeners and does not blockrefreshAsync()(including a refresh issued from the callback); the default callback and a throwing callback both leave the view working; a mapping failure during the initial replay no longer failsstart(); the reader is closed when the initial replay fails; a failedhasMessageAvailableAsync()failsstart()instead of hanging it; and the constructor rejects a configured compaction strategy without creating a reader.TableViewTest:testCreateMapped(custom mapping, value updates, tombstone vianullmapping),testCreateMappedWithIdentityMapper(TableView<Message<T>>with metadata access) andtestCreateMappedWithFailingMapper(a mapper that fails during the initial replay and while tailing; the key keeps its previous value and later messages are still applied)../gradlew quickCheckand the test classes above.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
Documentation
docdoc-requireddoc-not-neededdoc-completeMatching PR in forked repository
PR in forked repository: fork-repo