Skip to content

[improve][broker] Add metrics for message position find by timestamp - #26751

Merged
lhotari merged 3 commits into
apache:masterfrom
KannarFr:improve/message-find-metrics
Sep 29, 2026
Merged

lhotari merged 3 commits into
apache:masterfrom
KannarFr:improve/message-find-metrics

Conversation

@KannarFr

@KannarFr KannarFr commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

Seek / reset-cursor by timestamp and message TTL find a position with PersistentMessageFinder, which narrows the search to a range of ledgers using their close timestamps and then runs a binary search over entries (OpFindNewest), reading one entry per probe.

On ledgers offloaded to tiered storage each probe can be expensive: the offload index is sparse (one entry per data block, 64 MiB by default), so reading an entry in the middle of a block can require scanning the block from its start with ranged reads, each one costing tens of milliseconds on an object storage. A seek by timestamp into offloaded data can therefore take much longer than into BookKeeper while the subscription is fenced, and TTL checks on old topics go through the same path.

There is currently no way to observe this: no metric tells how long a find takes, how many entries it reads, or whether they come from BookKeeper or tiered storage. This makes it hard to diagnose slow seeks and to measure improvements to the search.

Modifications

  • Add OpenTelemetryMessageFinderStats with the following metrics. Attributes are intentionally limited to low-cardinality values (no topic or subscription), since a broker may own a very large number of topics:
    • pulsar.broker.message.find.duration (histogram, seconds, with bucket boundaries from 1 ms to 60 s), by pulsar.broker.message.find.reason (seek, expiry) and pulsar.broker.message.find.result (found, not_found, failure).
    • pulsar.broker.message.find.entry.read.count (counter, entries) and pulsar.broker.message.find.entry.read.size (counter, bytes), by reason and pulsar.broker.message.find.entry.storage (bookkeeper, offloaded).
  • The storage is the one the entry's ledger is actually read from: ManagedLedgerImpl.isReadFromOffloadedLedgerHandle(ledgerId) checks whether the opened read handle is an OffloadedLedgerHandle, so an offloaded ledger still read from BookKeeper (bookkeeper-first read priority) is counted as bookkeeper. The file system offloader's FileStoreBackedReadHandleImpl and BlobStoreBackedReadHandleImplV2 now implement the OffloadedLedgerHandle marker too; their other uses of the marker are unchanged since they keep the default lastAccessTimestamp() of -1, which opts out of the idle handle eviction.
  • The entry read metrics count the entries evaluated by the search: entries served from the broker entry cache are included, under the storage of their ledger. This is stated in the metric descriptions.
  • PersistentMessageFinder records the metrics for each entry read and when a find completes, and adds the per-find details (duration, entries and bytes read, offloaded entries read, whether the ledger range was narrowed) to the existing Found position closest to provided timestamp log line.
  • Wire the stats for the seek path (PersistentSubscription.resetCursor(long)) and the TTL paths (PersistentMessageExpiryMonitor, PersistentTopic). The previous constructor is kept for tests.

Follow-up: document the new metrics in the OpenTelemetry metrics reference on the website (apache/pulsar-site).

Verifying this change

  • Make sure that the change passes the CI checks.

This change added tests and can be verified as follows:

  • Added PersistentMessageFinderTest.testPersistentMessageFinderMetricsWhenFound (exact entry count and bytes, compared with the reads intercepted in BookKeeper), ...WhenNotFoundForExpiry, ...WhenReadFails (injected read failure) and ...WhenOffloaded.
  • Added FileSystemManagedLedgerOffloaderTest.testReadHandleIsOffloadedLedgerHandle and an assertion in BlobStoreManagedLedgerOffloaderStreamingTest.testReadAndWrite on the real read handles.
  • Existing tests PersistentMessageFinderTest and SubscriptionSeekTest pass.

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: new pulsar.broker.message.find.* OpenTelemetry metrics, see Modifications
  • Anything that affects deployment

🤖 Generated with Claude Code

Seek/reset-cursor by timestamp and message TTL find a position with a binary
search over entries. On offloaded ledgers each probe can be expensive, but
there is no way to observe how long a find takes or how many entries it reads.

Add OpenTelemetry metrics with low-cardinality attributes only (no topic or
subscription):
- pulsar.broker.message.find.duration, by reason (seek, expiry) and result
  (found, not_found, failure)
- pulsar.broker.message.find.entry.read.count and .size, by reason and by the
  storage the entry's ledger is read from (bookkeeper, offloaded)

Also log the per-find details (duration, entries and bytes read, offloaded
entries read, whether the ledger range was narrowed) on the existing
"Found position closest to provided timestamp" log line.

Assisted-by: Claude Code (Claude Opus 5.5)
Co-Authored-By: Claude Opus 5.5 <[email protected]>

@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 adding visibility into timestamp-based position searches; the low-cardinality attribute choice and the per-find completion accounting (found / not_found / failure, recorded once, concurrent-find rejections not counted) look right. Two things to address before merge: the storage attribute misclassifies several offload paths and cache hits, and the duration histogram has no seconds-scale bucket boundaries. I also suggest a few more test cases (see inline comments).

- Mark the file system and the V2 blob store read handles as
  OffloadedLedgerHandle, so that entries of ledgers offloaded with these
  handles are counted with storage=offloaded.
- Advise second-scale bucket boundaries (1 ms to 60 s) for the find
  duration histogram.
- Document in the entry read metrics that entries served from the broker
  entry cache are included, under the storage of their ledger.
- Test the not_found and failure results, the expiry reason, the
  offloaded storage, and the exact number of entries and bytes read.

Assisted-by: Claude Code (Claude Opus 5.5)
Co-Authored-By: Claude Opus 5.5 <[email protected]>

@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 quick follow-up. The offload markers, the cache documentation, the bucket advice and the extra test cases all look good. One test-stability point is left on the new tests (inline); it is a small change.

Open the durable cursor before writing the entries, so that its
mark-delete position keeps the ledgers from being trimmed in the
background before the find runs.

Assisted-by: Claude Code (Claude Opus 5.5)
Co-Authored-By: Claude Opus 5.5 <[email protected]>

@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. Thanks for the quick fix: opening the durable cursor before writing keeps the ledgers pinned in all three tests, and reopening it after the ledger is reopened restores the same cursor.

@lhotari lhotari added this to the 5.0.0 milestone Sep 29, 2026
@lhotari
lhotari merged commit b7dc4f9 into apache:master Sep 29, 2026
43 checks passed
@KannarFr
KannarFr deleted the improve/message-find-metrics branch September 30, 2026 23:37
lhotari pushed a commit to apache/pulsar-site that referenced this pull request Oct 1, 2026
…d by timestamp (#1236)

Document the pulsar.broker.message.find.* metrics added in apache/pulsar#26751:
the duration of a find by timestamp (subscription seek / reset-cursor, message
TTL), and the entries and bytes it reads, by storage (BookKeeper or tiered
storage).

Assisted-by: Claude Code (Claude Opus 5.5)

Co-authored-by: Claude Opus 5.5 <[email protected]>
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