Repository navigation
Rebase and retry the IdealState update in place on a version conflict during tier-relocation rebalance - #19069
Jackie-Jiang wants to merge 1 commit into
Conversation
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## master #19069 +/- ##
============================================
- Coverage 68.07% 68.05% -0.03%
Complexity 1450 1450
============================================
Files 3516 3516
Lines 228011 228057 +46
Branches 36077 36087 +10
============================================
- Hits 155213 155194 -19
- Misses 60613 60680 +67
+ Partials 12185 12183 -2
Flags with carried forward coverage won't be shown. Click here to find out more. ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
ddcba5e to
931f4af
Compare
There was a problem hiding this comment.
Pull request overview
This PR improves TableRebalancer’s behavior during strict-realtime tier-relocation rebalances by handling IdealState compare-and-set (ZK version) conflicts more efficiently. When a batch moves only tier segments and a concurrent IdealState write is disjoint from the segments in the batch, the rebalance now rebases the batch onto the latest IdealState and retries in place, avoiding an expensive restart of the convergence loop.
Changes:
- Add a bounded “rebase and retry” loop for
IdealStateupdates onZkBadVersionExceptionwhen the concurrent write does not touch the batch’s moved segments. - Refresh
targetAssignmentafter a successful rebase so the convergence check remains well-defined if segments were concurrently added/removed. - Introduce
MAX_IDEAL_STATE_UPDATE_REBASE_ATTEMPTSto bound in-place retries.
Comments suppressed due to low confidence (1)
pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/rebalance/TableRebalancer.java:918
- When the in-place rebase loop gives up (
updated == false), restoreexpectedVersionto the value associated with the still-currentcurrentAssignmentso the next convergence-loop iteration will reliably detect the IdealState version mismatch and refreshcurrentAssignment/targetAssignment. Without this, a fallback after exhausting rebase attempts can leaveexpectedVersionequal to the latest ZK version whilecurrentAssignmentis stale.
} else {
tableRebalanceLogger.info("Version changed while updating IdealState");
// Since IdealState wasn't updated, rollback the stats changes made
_tableRebalanceObserver.onRollback();
}
| boolean rebasable = isMovingOnlyTierSegments(segmentsToMove, providedTierToSegmentsMap); | ||
| boolean updated = false; | ||
| int rebaseAttempts = 0; | ||
| while (true) { |
There was a problem hiding this comment.
I think in this case it would go to https://github.com/apache/pinot/pull/19069/changes#diff-d4962fcb9ad5591bb650b16bec8858ae3d9aa357aadc5cf2940a0524216898aaR651, so it should be fine? Please confirm this as well.
| // Check version and update the IdealState. If the compare-and-set fails only because a concurrent write bumped | ||
| // the version without touching the segments this batch moves (e.g. consuming segment commits on a continuously | ||
| // ingesting table), rebase this batch onto the latest IdealState and retry the compare-and-set in place, without | ||
| // waiting for the ExternalView to converge again (this batch never landed, so there is nothing new to wait for) | ||
| // or recomputing the full target assignment. This keeps the rebalance from live-locking against a steady stream | ||
| // of version bumps. Only attempted when this rebalance moves only tier segments: the base placements are then | ||
| // unchanged, so a segment added concurrently keeps its correct placement and can be carried over as-is. |
J-HowHuang
left a comment
There was a problem hiding this comment.
Makes sense. I think we can probably reuse the code, see comments.
| // or recomputing the full target assignment. This keeps the rebalance from live-locking against a steady stream | ||
| // of version bumps. Only attempted when this rebalance moves only tier segments: the base placements are then | ||
| // unchanged, so a segment added concurrently keeps its correct placement and can be carried over as-is. | ||
| boolean rebasable = isMovingOnlyTierSegments(segmentsToMove, providedTierToSegmentsMap); |
There was a problem hiding this comment.
Should this be
| boolean rebasable = isMovingOnlyTierSegments(segmentsToMove, providedTierToSegmentsMap); | |
| boolean rebasable = !isStrictRealtimeSegmentAssignment || isMovingOnlyTierSegments(batchMovedSegments, providedTierToSegmentsMap); |
| for (String segment : batchMovedSegments) { | ||
| rebasedAssignment.put(segment, nextAssignment.get(segment)); | ||
| } | ||
| nextAssignment = rebasedAssignment; |
There was a problem hiding this comment.
nit: Can we have a method rebaseTargetAssignment(originalTargetAssignment, originalCurrentAssignment, newCurrentAssingment) that resolves rebasable and rebasedAssignment?
Then we can use the same method here: https://github.com/apache/pinot/pull/19069/changes#diff-d4962fcb9ad5591bb650b16bec8858ae3d9aa357aadc5cf2940a0524216898aaR651-R720, it's basically doing the same thing only that it rebases target assignment and here we rebase next assignment
| Map<String, Map<String, String>> refreshedTarget = new TreeMap<>(currentAssignment); | ||
| for (String segment : SegmentAssignmentUtils.getSegmentsToMove(currentAssignment, targetAssignment)) { | ||
| if (currentAssignment.containsKey(segment)) { | ||
| refreshedTarget.put(segment, targetAssignment.get(segment)); | ||
| } | ||
| } | ||
| targetAssignment = refreshedTarget; |
There was a problem hiding this comment.
we may also reuse the rebaseTargetAssignment method if we have one here.
| boolean rebasable = isMovingOnlyTierSegments(segmentsToMove, providedTierToSegmentsMap); | ||
| boolean updated = false; | ||
| int rebaseAttempts = 0; | ||
| while (true) { |
There was a problem hiding this comment.
I think in this case it would go to https://github.com/apache/pinot/pull/19069/changes#diff-d4962fcb9ad5591bb650b16bec8858ae3d9aa357aadc5cf2940a0524216898aaR651, so it should be fine? Please confirm this as well.
…ier-relocation rebalance When TableRebalancer moves segments in batches, each batch is applied with a version-checked IdealState update. On a continuously-ingesting strict realtime table (upsert/dedup), consuming-segment commits bump the IdealState version between the read and the write, so the update loses the compare-and-set. Previously every lost compare-and-set fell back to the top of the convergence loop, paying another ExternalView-convergence wait and target recompute, so the rebalance could make very little progress while ingestion continued. When the rebalance moves only tier segments (so the base placements of the partitions are unchanged), a concurrent write that does not touch the segments this batch moves cannot invalidate the batch. In that case, on a version conflict, re-read the IdealState, and if the concurrent change is disjoint from the segments this batch moves, rebase the batch onto the latest IdealState and retry the compare-and-set in place, without waiting for the ExternalView to converge again or recomputing the target. The retry is bounded; on overlap, re-read failure, or exhausting the retries, it falls back to the convergence loop as before. This builds on the target-recompute skip for tier-relocation rebalances: it reuses the same "moves only tier segments" condition to decide when a batch is safe to rebase, and refreshes the target from the adopted IdealState so the convergence check stays well-defined when a rebase pulls in concurrently added or removed segments.
931f4af to
47762a4
Compare
PR flow
IdealState update loop with rebase on version conflict for tier-only moves
AI-generated · Green: added · Yellow: modified · Red: removed · Gray: existing
Diff evidence
Summary
Follow-up to #19054. That PR reduces the cost of the version-checked
IdealStateupdate losing the compare-and-set during a tier-relocation rebalance of a strict realtime table (upsert/dedup); this PR reduces the cost of recovering from a lost compare-and-set.Each batch of segment moves is applied with a version-checked
IdealStateupdate. On a continuously-ingesting table, consuming-segment commits bump theIdealStateversion between the read and the write, so the update loses the compare-and-set. Previously every lost compare-and-set fell back to the top of the convergence loop — another ExternalView-convergence wait (~hundreds of ms) plus a target recompute — so the rebalance could make very little progress while ingestion continued.Change
When the rebalance moves only tier segments (the base placements of the partitions are unchanged — the same condition #19054 uses to skip the target recompute), a concurrent write that does not touch the segments this batch moves cannot invalidate the batch. So on a
ZkBadVersionException:IdealState.IdealState(reapply the batch's moves, preserving the concurrent changes) and retry the compare-and-set in place — no ExternalView wait, norebalanceTable.MAX_IDEAL_STATE_UPDATE_REBASE_ATTEMPTS), fall back to the convergence loop as before.On a successful rebase, the target assignment is refreshed from the adopted
IdealStateso thecurrentAssignment.equals(targetAssignment)convergence check stays well-defined when a rebase pulls in concurrently added (e.g. new consuming) or removed (e.g. retention) segments.Net effect: disjoint version churn is absorbed in place, so a tier relocation makes progress against a steady stream of consuming-segment commits instead of repeatedly restarting.
Notes
onRollbackis skipped since the batch did land) is the subtlest part and the best target for a unit test that injects a version bump between the read and the compare-and-set.