yashmayya opened a new pull request, #19643:
URL: https://github.com/apache/pinot/pull/19643
## Problem
`TimeSegmentPruner.refreshSegment()` rebuilds the table's entire interval
tree every time one segment's ZK metadata is refreshed. `_intervalMap` holds
every segment of the table, and the `IntervalTree` constructor walks all of it,
so a single-segment refresh costs O(n) allocation and O(n log n) CPU over all
segments, inside a `synchronized` method.
This is not limited to explicit segment refreshes. Every LLC REALTIME
segment commit goes through
`PinotLLCRealtimeSegmentManager.updateSegmentZKMetadataToDone()`, which calls
`sendSegmentRefreshMessage(realtimeTableName, segmentName, false, true)`. That
message is **broadcast to every broker**:
```
PinotLLCRealtimeSegmentManager.updateSegmentZKMetadataToDone(...)
controller
-> sendSegmentRefreshMessage(table, segment, false, true)
[Helix USER_DEFINE_MSG -> every broker]
BrokerUserDefinedMessageHandlerFactory$RefreshSegmentMessageHandler.handleMessage()
BaseBrokerRoutingManager.refreshSegment(String, String)
SegmentZkMetadataFetcher.refreshSegment(String)
TimeSegmentPruner.refreshSegment(String, ZNRecord)
new IntervalTree<>(_intervalMap) //
full rebuild, ALL segments
```
So the per-broker rebuild rate is the cluster-wide commit rate, while the
per-broker query rate is the cluster query rate divided by the broker count. On
a large REALTIME table this makes broker heap churn quadratic in segment count
for a fixed commit rate.
This was found on a production broker with a 487k-segment table (253k
REALTIME). A JFR allocation profile attributed 11.5% of sampled allocations to
`IntervalTree.<init>` under `TimeSegmentPruner.refreshSegment`, the largest
single identified Pinot allocation site. The broker's live set was small (~3
GiB against a 29 GiB `Xmx`), but old-gen reclamation ran less than once per 5
minutes, so this garbage was tenured rather than dying in eden.
## Change
Two commits, separately bisectable.
**1. Build `IntervalTree` from flat arrays.** The tree now holds the sorted
distinct intervals, a subtree-max array, a value array and an offset array. The
balanced tree over them is implicit: the root of index range `[start, end)` is
`start + (end - start) / 2`. This is the same tree shape the old code
materialized, so search semantics are unchanged. Removed: the intermediate
`HashMap<Interval, List<VALUE>>`, one `ArrayList` per distinct interval, one
`Node` per distinct interval, and the `LinkedList<Pairs.IntPair>` used for the
breadth-first layout walk.
**2. Rebuild the tree lazily on read.** A segment change sets
`_intervalTree` to `null`; `prune()` rebuilds it under double-checked locking,
and only after it knows the query has a prunable time filter. Two invalidations
are also skipped: `refreshSegment` when the new interval equals the old one
(the common OFFLINE re-push), and `onAssignmentChange` when no segment was
added or removed.
Rebuilding **on read** rather than on a timer is load-bearing, not just
convenient. `prune()` takes the segments it returns *from* the tree and uses
its `segments` argument only as a filter, so a segment missing from the tree is
dropped from the routing rather than merely left unpruned. Any scheme that lets
a query observe a stale tree silently loses data. Rebuilding on read means the
tree is always complete at the moment it is read.
## Benchmark results
New JMH benchmark `BenchmarkTimeSegmentPruner`. The numbers below come from
a standalone harness using
`com.sun.management.ThreadMXBean.getThreadAllocatedBytes`, which counts exact
bytes rather than sampling. JDK 25, `-Xmx8g`, same harness against both
revisions.
250,000 segments, one distinct interval per segment (REALTIME time ranges
come from the data in milliseconds, so intervals are effectively all distinct):
| Workload | master | this PR | change |
|---|---|---|---|
| 1 segment refresh | 70.47 MB, 84 ms | 0.0005 MB, <0.01 ms | work is
deferred |
| 1 refresh + 1 query | 70.47 MB, 99.6 ms | 7.63 MB, 31.9 ms | 9.2x alloc,
3.1x CPU |
| 6 refreshes + 1 query | 422.8 MB, 620 ms | 7.64 MB, 32.9 ms | 55x alloc,
19x CPU |
| 1 query, tree already built | 2.7 KB, ~1 us | 2.1 KB, ~1 us | unchanged |
The 6:1 row is the ratio implied by broadcasting refreshes to 6 brokers
while splitting queries across them.
Other sizes and shapes (`dup` = segments sharing one interval; `dup=64`
models a coarse time unit such as DAYS):
| Segments | dup | Workload | master | this PR |
|---|---|---|---|---|
| 10,000 | 1 | 6 refreshes + 1 query | 17.25 MB, 7.0 ms | 0.31 MB, 0.83 ms |
| 100,000 | 1 | 6 refreshes + 1 query | 170.4 MB, 155 ms | 3.06 MB, 11.5 ms |
| 250,000 | 1 | 6 refreshes + 1 query | 422.8 MB, 620 ms | 7.64 MB, 32.9 ms |
| 250,000 | 64 | 6 refreshes + 1 query | 60.1 MB, 44.9 ms | 3.98 MB, 24.2 ms
|
| 250,000 | 64 | 1 refresh + 1 query | 10.5 MB, 8.9 ms | 3.98 MB, 23.4 ms |
Scaled to the rates measured on the production broker (1,763 refreshes/hour
and at most 1,320 time-filtered queries/hour reaching one broker for that
table): roughly 121 GB/hour of garbage on master against at most 9.8 GB/hour
here.
Two independent effects combine. The cheaper rebuild is a flat ~9.2x
whatever the rates are. Deferral cuts the rebuild *count* from one per change
to roughly `min(changes, time-filtered queries)`. With `C` the per-broker
change rate and `Q` the per-broker rate of queries that read the tree, the
allocation win is about `9.2 x max(1, C/Q)`. It never falls below 9.2x; the
ratio only sets the ceiling.
## Trade-offs
**The rebuild moves onto the query path.** The first query with a prunable
time filter after a segment change builds the tree — 32 ms at 250,000 segments
— and queries arriving during that build block on the pruner monitor. When the
tree is already built, `prune()` is one volatile read and takes no lock,
unchanged from master. The work it replaces was 84 ms per refresh on the Helix
message-handler thread, so total work drops, but the latency lands somewhere
new. Queries without a prunable time filter never trigger a rebuild.
**One CPU regression, in the grouped-interval shape.** The new constructor
sorts all values rather than grouping by interval first and sorting only the
distinct ones. When many segments share an interval *and* queries outnumber
segment changes, a rebuild costs about 2.6x more CPU (23 ms against 9 ms at
250k/`dup=64`), while still allocating 2.6x less. Grouping first is not a free
fix — it needs a hash structure over the values. Measured over 250k values:
| | every interval distinct | 64 values per interval |
|---|---|---|
| Sort all values (this PR) | 19.5 ms, 2.86 MB | 11.8 ms, 2.86 MB |
| Group first, then sort | 35.1 ms, 11.91 MB | 3.1 ms, 1.13 MB |
Grouping first costs the motivating case 15 ms and 9 MB, so this keeps a
single code path. A comment in the constructor records the measurement.
**A failure mode changed from loud to silent.** `prune()` before `init()`
used to throw NPE and now returns an empty segment set. This is unreachable
today: `BaseBrokerRoutingManager` calls `init` before publishing the routing
entry.
## Caveats
- No config flag. The behaviour change applies to every table.
- Below roughly 100,000 segments none of this is material — master already
only spends 2.9 MB and 1.2 ms per rebuild at 10,000 segments.
- This is a lazy whole-tree rebuild, not an incremental tree. Removing the
query-path rebuild entirely needs either a genuinely incremental structure, or
a pending-set overlay that lets `prune()` add back segments changed since the
last build. Both are larger changes and are left as follow-ups.
- `SegmentPartitionMetadataManager.refreshSegment()` has the same shape — it
calls `computeAllTablePartitionInfo()`, which walks every segment of the table
per refresh, and accounted for 3.87% of the same profile. Not addressed here.
`TimeBoundaryManager.refreshSegment()` rescans all segment end times per
refresh but allocates nothing.
## Testing
- `IntervalTreeTest`: 4 new cases — empty tree, null search interval, all
values on one interval, and a randomized cross-check of ~3,500 searches against
a brute force scan over 35 tree shapes.
- `SegmentPrunerTest`: 2 new cases — repeated refresh (covers the
unchanged-interval short circuit with `assertSame`/`assertNotSame` on the tree,
and several refreshes landing between two queries), and concurrent segment
change against concurrent pruning.
- Negative control verified: removing the `synchronized` block from
`getIntervalTree()` makes the concurrency test fail with
`ConcurrentModificationException`, surfaced through `Future.get()`.
- Offline differential check against the previous implementation: 396
randomly generated trees and 15,840 searches, where the new tree, the old tree
and a brute force scan all agree.
- All 219 tests in `org.apache.pinot.broker.routing` pass.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]