Skip to content

[feat][client] PIP-445: Add Builder Methods to Create Message-based TableView - #24809

Merged
lhotari merged 17 commits into
apache:masterfrom
namest504:master
Sep 21, 2026
Merged

lhotari merged 17 commits into
apache:masterfrom
namest504:master

Conversation

@namest504

@namest504 namest504 commented Oct 2, 2025 •

Copy link
Copy Markdown
Contributor

Fixes: #24744

Main Issue: #24744

PIP: 445 #24842

Motivation

The current TableView API 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 adds createMapped methods to the TableViewBuilder, allowing users to transform a full Message<T> into any custom object V that 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

  • Added TableViewMessageMapper<T, V> to the public client API: a functional interface with V map(Message<T>), which may throw, and a default boolean onMappingError(Message<T>, Throwable) callback. The default callback returns false, which makes the table view log the failure at ERROR level with the topic, key and message id; an implementation can override it to report the failure its own way (metrics, signalling another thread, ...) and return true to suppress the log. Further callbacks can be added later as default methods without breaking implementations.
  • Added createMapped and createMappedAsync to TableViewBuilder: they take a TableViewMessageMapper<T, V>. Lambdas work directly, e.g. createMapped(msg -> msg.getProperty("region")), and createMapped(msg -> msg) builds a TableView<Message<T>> exposing the complete messages. A null mapper is rejected with IllegalArgumentException (failed future in the async variant).
  • Defined the mapping contract on the public API: a keyed message with an empty payload is a tombstone that removes the key without calling the mapper; a null return 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 pending refresh() still completes — and onMappingError is called. Exceptions thrown by the callback are logged and ignored.
  • Rejected topicCompactionStrategyClassName for mapped views in both builder methods and, defensively, in the MappedTableViewImpl constructor before the reader is created: a TopicCompactionStrategy compares values of the topic's schema type, which a mapped view does not store, so the combination would fail with a ClassCastException on the reader thread.
  • Moved the core logic to AbstractTableViewImpl<T, V>: TableViewImpl keeps 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 failing create() in the former and being dropped silently in the latter. This also applies to the broker's own create() user, the extensible load balancer's ServiceUnitStateTableViewImpl: an undecodable ServiceUnitStateData record no longer prevents the channel from starting; it is skipped and logged at ERROR, as it already was while tailing. The trade-off is documented in the PIP-445 compatibility section.
  • Introduced package-private MappedTableViewImpl<T, V>: implements the mapped views and delegates mapping failures to the mapper's onMappingError. Message pooling is disabled for mapped views since the mapper may retain the Message instance; the createMapped javadoc notes that a retained Message also keeps its metadata, schema and connection reference, so a value object is preferable to msg -> msg for topics with many keys.
  • Closed the reader when start() fails: the reader created for the table view is no longer leaked when the initial replay fails. A failed hasMessageAvailableAsync() during the replay now fails start() too, instead of leaving create() waiting forever with the reader open.

Verifying this change

  • Unit tests in TableViewBuilderImplTest: testCreateMapped* (mapped view creation, null mapper validation, and rejection of a configured compaction strategy for both the sync and async variants).
  • New unit tests in MappedTableViewImplTest with a mocked reader: a mapping failure skips the message and keeps the previous value, notifies onMappingError, does not notify listeners and does not block refreshAsync() (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 fails start(); the reader is closed when the initial replay fails; a failed hasMessageAvailableAsync() fails start() instead of hanging it; and the constructor rejects a configured compaction strategy without creating a reader.
  • Integration tests in TableViewTest: testCreateMapped (custom mapping, value updates, tombstone via null mapping), testCreateMappedWithIdentityMapper (TableView<Message<T>> with metadata access) and testCreateMappedWithFailingMapper (a mapper that fails during the initial replay and while tailing; the key keeps its previous value and later messages are still applied).
  • Verified locally with ./gradlew quickCheck and the test classes above.

Does this pull request potentially affect one of the following parts:

If the box was checked, please highlight the changes

  • Dependencies (add or upgrade a dependency)
  • The public API
  • The schema
  • The default values of configurations
  • The threading model
  • The binary protocol
  • The REST endpoints
  • The admin CLI options
  • The metrics
  • Anything that affects deployment

Documentation

  • doc
  • doc-required
  • doc-not-needed
  • doc-complete

Matching PR in forked repository

PR in forked repository: fork-repo

@github-actions github-actions Bot added the doc-required Your PR changes impact docs and you will update later. label Oct 2, 2025
…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.

@lhotari lhotari left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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).

Comment thread pulsar-client/src/main/java/org/apache/pulsar/client/impl/TableViewImpl.java Outdated
@namest504

Copy link
Copy Markdown
Contributor Author

Hi @lhotari,

Thanks for the review and the great suggestion about using the TableViewBuilder!

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.

@lhotari

lhotari commented Oct 14, 2025 •

Copy link
Copy Markdown
Member

Hi @lhotari,

Thanks for the review and the great suggestion about using the TableViewBuilder!

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).

@namest504 namest504 changed the title [improve][api] Add getRawMessage() method to TableView for accessing raw Pulsar message [improve][api] Add Builder Methods to Create Message-based TableView Oct 20, 2025
@namest504
namest504 requested a review from lhotari October 20, 2025 04:47
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.
@namest504

Copy link
Copy Markdown
Contributor Author

@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:

  • Kept the original TableViewImpl class name and moved the shared logic to a package-private AbstractTableViewImpl<T, V>, following the naming in your experiment branch. No imports are removed and there are no fully qualified class name references anymore.
  • Removed the unused MessageTableView class left over from the earlier createForMessages approach.
  • Disabled message pooling for mapped views: the mapper may retain the Message instance, so the pooled payload buffer must not be reused or released. The classic TableView path keeps pooling and releasing as before.
  • Added mapper null-validation and unit tests in TableViewBuilderImplTest.

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.

@namest504 namest504 changed the title [improve][api] Add Builder Methods to Create Message-based TableView [improve][client] Add Builder Methods to Create Message-based TableView Jul 10, 2026
…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.
@namest504
namest504 requested a review from lhotari July 10, 2026 03:23
lhotari and others added 2 commits August 31, 2026 12:36
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 lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@lhotari

lhotari commented Sep 16, 2026

Copy link
Copy Markdown
Member

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 .poolMessages(true) with its unconditional release(), while the mapped path opts out, so nothing can release a buffer your mapper still holds.

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 createMapped is new public API and that behaviour becomes fixed once released. One of the behaviours I flag is pre-existing on master and I have said so in that comment; the concern is that a mapper is arbitrary user code where a schema-decode failure was not. Congratulations on the PIP-445 vote.

…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
@github-actions github-actions Bot added the PIP label Sep 17, 2026
# Conflicts:
#	pulsar-client/src/main/java/org/apache/pulsar/client/impl/TableViewImpl.java
@lhotari

lhotari commented Sep 17, 2026

Copy link
Copy Markdown
Member

@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: createMapped takes a TableViewMessageMapper (map may throw, onMappingError reports failures), mapped views reject compaction strategies, the impl is the package-private MappedTableViewImpl, tombstones are documented, and PIP-445 plus the PR description are updated to match. I also merged master to resolve the conflict with #26626 and ported that fix into AbstractTableViewImpl. Each review thread has a short note on how it was handled; please take a look and let me know if anything looks off.

@lhotari lhotari left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@lhotari lhotari changed the title [improve][client] Add Builder Methods to Create Message-based TableView [feat][client] PIP-445: Add Builder Methods to Create Message-based TableView Sep 17, 2026
@lhotari lhotari added this to the 5.0.0 milestone Sep 17, 2026
@david-streamlio

Copy link
Copy Markdown
Contributor

Nice work — the refactor is a faithful extraction. I diffed AbstractTableViewImpl line-by-line against the base TableViewImpl at a846a45953 and the only core-logic deltas are the ones the PIP claims: generify to <T, V>, the poolMessages flag, close-on-start-failure, and the guarded value extraction. createMappedAsync correctly returns failed futures rather than throwing synchronously (TableViewBuilderImpl.java:88-93), there's no new reflection into private state, no blocking on the reader thread, and no new Recycler usage.

I ran ./gradlew :pulsar-client-original:test --tests MappedTableViewImplTest --tests TableViewBuilderImplTest -PtestRetryCount=0 → 22/22 green. I did not run the pulsar-broker tests (TableViewTest, ServiceUnitStateCompactionTest).

Everything below is non-blocking; the PR already has an approval. Two points I'd like your read on, then some nits.

1. The create() behavior change lands on the broker's own load-balancer path

AbstractTableViewImpl.java:249-262 now swallows a decode failure for plain table views too, not just mapped ones. ServiceUnitStateTableViewImpl.java:94 builds its channel view with create(), so a ServiceUnitStateData message that fails to decode used to fail channel start loudly and now leaves a silent hole in the ownership view, observable only as an ERROR log.

The PIP documents the general change in its compatibility section, but doesn't name this consumer. Is that trade deliberate? Skipping is probably the better operational default — a broker that won't start its load-balancer channel is worse than one that drops a corrupt record — but create() users have no programmatic signal at all. Worth either calling the consumer out explicitly in the PIP, or considering whether non-mapped views should get an equivalent of onMappingError at some point.

2. createMapped(msg -> msg) retains more than the message

The pooling claim itself checks out: with poolMessage == false, MessageImpl.java:207 does Unpooled.copiedBuffer(payload), so retaining the Message is safe. But MessageImpl.java:196 also assigns msg.cnx = cnx, and the instance holds msgMetadata, brokerEntryMetadata and schema. So an identity-mapped view over a large keyspace pins a ClientCnx reference per stored key, keeping old connection objects reachable across reconnects.

Safe, just not cheap. A sentence in the TableViewBuilder#createMapped javadoc recommending extraction of the needed fields over msg -> msg for high-cardinality topics would set the right expectation.

Nits

  • AbstractTableViewImpl.java:135-149 — the reader-leak fix doesn't cover the adjacent path. readAllExistingMessages (line 434) has no .exceptionally on hasMessageAvailableAsync(), so if that call fails, future never completes, start() hangs, and the new whenComplete never fires — reader still leaked. Pre-existing, but you're fixing the leak right next to it; one .exceptionally(ex -> { future.completeExceptionally(ex); return null; }) would close the case.
  • AbstractTableViewImpl.java:251 — catch (Throwable t) routes OutOfMemoryError / StackOverflowError into onMappingError and keeps reading. map declares throws Exception, so catch (Exception e) matches the contract and lets a VirtualMachineError escape.
  • AbstractTableViewImpl.java:259 — "Failed to map message, skipping it" reads oddly for plain TableViewImpl, where there is no mapper and the cause is a decode failure.
  • TableViewBuilderImpl.java:77-78 — the two checkArgument calls are redundant: PulsarClientException.unwrap rethrows RuntimeException unchanged (PulsarClientException.java:1058-1059), so the async guards already surface IllegalArgumentException from createMapped(). Harmless, and it does fail fast before a reader is created; just two places to keep in sync.
  • MappedTableViewImpl.java:37-44 — the compaction-strategy invariant is enforced only in the builder. A defensive check in the constructor would keep it local to the class that would otherwise hit a ClassCastException on the reader thread.
  • Making AbstractTableViewImpl package-private while TableViewImpl stays public is what forced readDeclaredField → readField in TableViewTest and ServiceUnitStateCompactionTest. Fine inside client.impl, just noting the reflective-access consequence is real for anyone doing the same downstream.

Adding two abstract methods to TableViewBuilder looks fine to me — it's @InterfaceStability.Evolving and TableViewBuilderImpl is the only implementation in-tree.

…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
@lhotari

lhotari commented Sep 18, 2026

Copy link
Copy Markdown
Member

@david-streamlio thanks for the thorough second pass — the line-by-line diff against the base TableViewImpl is exactly the kind of check this extraction needed. Here is what I did with each point.

1. create() behaviour change on the load balancer channel. Deliberate, and now called out by name. The PIP-445 compatibility section names ServiceUnitStateTableViewImpl and spells out the trade: on master a ServiceUnitStateData that fails to decode was already skipped (with a log that blamed the reader) when it arrived while tailing; only the replay phase changes, from "channel does not start" to "record skipped, logged at ERROR with topic, key and message id". A broker that cannot start its load balancer channel because of one undecodable record is the worse outcome. A programmatic failure signal for plain create() views (an onMappingError equivalent) is noted as follow-up material rather than added here — it would be new public API on the classic path.

2. msg -> msg retains more than the payload. Correct — MessageImpl.create assigns cnx, and the instance also carries msgMetadata, brokerEntryMetadata and the schema. The createMapped javadoc now says so and recommends copying the needed fields into a value object for topics with many keys.

Nits

  • hasMessageAvailableAsync() failure hangs create() — fixed. The outer chain in readAllExistingMessages now fails the replay future, so start() completes exceptionally and the whenComplete closes the reader. New test testStartFailsWhenHasMessageAvailableFails; with the .exceptionally removed it fails with a TimeoutException after 5 s, which is the hang.
  • catch (Throwable) around the mapper — kept, on purpose. The mapper runs inside CompletableFuture callbacks, and CompletableFuture catches Throwable regardless, so an Error cannot escape to the JVM either way: the only choice is whether it is attributed to the mapper (with key and message id, via onMappingError) or misdiagnosed as a reader failure and retried — which is the misattribution the guard exists to remove. StackOverflowError from a recursive mapper is the realistic case, and by the time it is caught the stack has unwound. It also matches how the client treats every other user callback: ConsumerBase catches Throwable around MessageListener.received, and the table view listener loop has always caught Throwable.
  • Log wording — changed to "Skipping message whose value could not be decoded or mapped", which reads correctly for both the classic and the mapped path.
  • Redundant checkArgument in createMapped — they are load-bearing, not redundant. unwrap rethrows a RuntimeException only when it is the top-level throwable; .get() throws a checked ExecutionException, and in the cause-unwrapping branch there is no IllegalArgumentException case, so it falls through to the generic new PulsarClientException(t). I verified by deleting the sync null check: testCreateMappedWhenMapperIsNull then fails with Expected exception of type IllegalArgumentException but got PulsarClientException: java.util.concurrent.ExecutionException: java.lang.IllegalArgumentException: mapper cannot be null. Without them createMapped() would break the documented @throws IllegalArgumentException contract.
  • Compaction-strategy invariant local to MappedTableViewImpl — done. The constructor rejects a configured strategy itself, via a static helper evaluated in the super(...) argument list so it runs before the reader is created (the client targets release 17, so no flexible constructor bodies). The class now owns the message constant and the builder reuses it. New test testConstructorRejectsCompactionStrategyBeforeCreatingReader asserts the IllegalArgumentException and verify(client, never()).newReader(...).
  • readDeclaredField → readField — agreed it is a real consequence for downstream reflective access; it stays as is since AbstractTableViewImpl is package-private on purpose.

@namest504

Copy link
Copy Markdown
Contributor Author

@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 onMappingError means the skip is not silent. Having that spelled out in the javadoc as a contract is something I would not have arrived at on my own. I went through the changes and the follow-ups from @david-streamlio's review as well, and nothing looks off to me.

@lhotari
lhotari merged commit 1b5ec96 into apache:master Sep 21, 2026
45 checks passed
Radiancebobo pushed a commit to Radiancebobo/pulsar that referenced this pull request Oct 8, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

doc-required Your PR changes impact docs and you will update later. PIP

Projects

None yet

Development

Successfully merging this pull request may close these issues.

API to allow access to raw message for TableView

3 participants