Repository navigation
[improve][broker] Add metrics for message position find by timestamp - #26751
Conversation
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
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
…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]>
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
OpenTelemetryMessageFinderStatswith 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), bypulsar.broker.message.find.reason(seek,expiry) andpulsar.broker.message.find.result(found,not_found,failure).pulsar.broker.message.find.entry.read.count(counter, entries) andpulsar.broker.message.find.entry.read.size(counter, bytes), by reason andpulsar.broker.message.find.entry.storage(bookkeeper,offloaded).ManagedLedgerImpl.isReadFromOffloadedLedgerHandle(ledgerId)checks whether the opened read handle is anOffloadedLedgerHandle, so an offloaded ledger still read from BookKeeper (bookkeeper-firstread priority) is counted asbookkeeper. The file system offloader'sFileStoreBackedReadHandleImplandBlobStoreBackedReadHandleImplV2now implement theOffloadedLedgerHandlemarker too; their other uses of the marker are unchanged since they keep the defaultlastAccessTimestamp()of -1, which opts out of the idle handle eviction.PersistentMessageFinderrecords 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 existingFound position closest to provided timestamplog line.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
This change added tests and can be verified as follows:
PersistentMessageFinderTest.testPersistentMessageFinderMetricsWhenFound(exact entry count and bytes, compared with the reads intercepted in BookKeeper),...WhenNotFoundForExpiry,...WhenReadFails(injected read failure) and...WhenOffloaded.FileSystemManagedLedgerOffloaderTest.testReadHandleIsOffloadedLedgerHandleand an assertion inBlobStoreManagedLedgerOffloaderStreamingTest.testReadAndWriteon the real read handles.PersistentMessageFinderTestandSubscriptionSeekTestpass.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
pulsar.broker.message.find.*OpenTelemetry metrics, see Modifications🤖 Generated with Claude Code