Skip to content
Draft
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,10 @@ public class TableRebalancer {
// TODO: Consider making the timeoutMs below table rebalancer configurable
private static final int TABLE_SIZE_READER_TIMEOUT_MS = 30_000;
private static final int STREAM_PARTITION_OFFSET_READ_TIMEOUT_MS = 10_000;
// Max number of times to rebase a single IdealState batch update onto a newer version and retry the compare-and-set
// in place (without waiting for the ExternalView to converge again) when the version changed only due to concurrent
// writes that do not touch the segments being moved. Beyond this, fall back to the regular convergence loop.
private static final int MAX_IDEAL_STATE_UPDATE_REBASE_ATTEMPTS = 5;
private static final AtomicInteger REBALANCE_JOB_COUNTER = new AtomicInteger(0);
private final HelixManager _helixManager;
private final HelixDataAccessor _helixDataAccessor;
Expand Down Expand Up @@ -801,27 +805,105 @@ private RebalanceResult doRebalance(TableConfig tableConfig, RebalanceConfig reb
idealState.setNumPartitions(nextAssignment.size());
idealState.setReplicas(Integer.toString(nextAssignment.values().iterator().next().size()));

// Check version and update IdealState
try {
Preconditions.checkState(_helixDataAccessor.getBaseDataAccessor()
.set(idealStatePropertyKey.getPath(), idealStateRecord, expectedVersion, AccessOption.PERSISTENT),
"Failed to update IdealState");
// Segments this batch changes relative to the current assignment it was computed against. Captured before the
// compare-and-set so that, on a version conflict, we can tell whether a concurrent write touched any segment
// this batch moves.
List<String> batchMovedSegments = SegmentAssignmentUtils.getSegmentsToMove(currentAssignment, nextAssignment);

// 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.
Comment on lines +813 to +819
boolean rebasable = isMovingOnlyTierSegments(segmentsToMove, providedTierToSegmentsMap);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

Should this be

Suggested change
boolean rebasable = isMovingOnlyTierSegments(segmentsToMove, providedTierToSegmentsMap);
boolean rebasable = !isStrictRealtimeSegmentAssignment || isMovingOnlyTierSegments(batchMovedSegments, providedTierToSegmentsMap);

boolean updated = false;
int rebaseAttempts = 0;
while (true) {
try {
Preconditions.checkState(_helixDataAccessor.getBaseDataAccessor()
.set(idealStatePropertyKey.getPath(), idealStateRecord, expectedVersion, AccessOption.PERSISTENT),
"Failed to update IdealState");
updated = true;
break;
} catch (ZkBadVersionException e) {
if (!rebasable || rebaseAttempts >= MAX_IDEAL_STATE_UPDATE_REBASE_ATTEMPTS) {
break;
}
IdealState latestIdealState;
try {
latestIdealState = _helixDataAccessor.getProperty(idealStatePropertyKey);
} catch (Exception re) {
tableRebalanceLogger.warn("Failed to re-read IdealState for rebasing after a version conflict", re);
break;
}
if (latestIdealState == null) {
break;
}
Map<String, Map<String, String>> latestAssignment = latestIdealState.getRecord().getMapFields();
boolean concurrentChangeTouchesBatch = false;
for (String segment : batchMovedSegments) {
if (!Objects.equals(currentAssignment.get(segment), latestAssignment.get(segment))) {
concurrentChangeTouchesBatch = true;
break;
}
}
if (concurrentChangeTouchesBatch) {
break;
}
// Rebase the batch onto the latest IdealState, preserving concurrent changes to the other segments.
Map<String, Map<String, String>> rebasedAssignment = new TreeMap<>(latestAssignment);
for (String segment : batchMovedSegments) {
rebasedAssignment.put(segment, nextAssignment.get(segment));
}
nextAssignment = rebasedAssignment;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

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

idealState = latestIdealState;
idealStateRecord = latestIdealState.getRecord();
idealStateRecord.setMapFields(nextAssignment);
idealState.setNumPartitions(nextAssignment.size());
idealState.setReplicas(Integer.toString(nextAssignment.values().iterator().next().size()));
expectedVersion = idealStateRecord.getVersion();
rebaseAttempts++;
tableRebalanceLogger.info("Rebasing IdealState update onto version {} after a concurrent change that does "
+ "not affect the segments being moved (attempt {})", expectedVersion, rebaseAttempts);
} catch (Exception e) {
onReturnFailure("Caught exception while updating IdealState, aborting the rebalance", e,
tableRebalanceLogger);
return new RebalanceResult(rebalanceJobId, RebalanceResult.Status.FAILED,
"Caught exception while updating IdealState: " + e, instancePartitionsMap, tierToInstancePartitionsMap,
targetAssignment, preChecksResult, summaryResult);
}
}

if (updated) {
expectedVersion++;
currentAssignment = nextAssignment;
if (rebaseAttempts > 0) {
// A rebase adopts the latest IdealState, which may have segments added (e.g. new consuming segments) or
// removed (e.g. by retention) concurrently that the target assignment, computed earlier, does not match.
// Since a rebase is only done when the rebalance moves only tier segments, added segments keep their current
// placement, so rebuild the target from the adopted current assignment while preserving the target only for
// the segments still moving that still exist. This keeps the convergence check
// (currentAssignment.equals(targetAssignment)) well-defined and its key set aligned with the current
// 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;
Comment on lines +890 to +896

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

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

we may also reuse the rebaseTargetAssignment method if we have one here.

}
segmentsToMove = SegmentAssignmentUtils.getSegmentsToMove(currentAssignment, targetAssignment);
// IdealState update is successful. Update the segment list as the IDEAL_STATE_CHANGE_TRIGGER should have
// captured the newly added / deleted segments
allSegmentsFromIdealState = currentAssignment.keySet();
tableRebalanceLogger.info("Successfully updated the IdealState");
} catch (ZkBadVersionException e) {
} else {
tableRebalanceLogger.info("Version changed while updating IdealState");
// Since IdealState wasn't updated, rollback the stats changes made
_tableRebalanceObserver.onRollback();
} catch (Exception e) {
onReturnFailure("Caught exception while updating IdealState, aborting the rebalance", e, tableRebalanceLogger);
return new RebalanceResult(rebalanceJobId, RebalanceResult.Status.FAILED,
"Caught exception while updating IdealState: " + e, instancePartitionsMap, tierToInstancePartitionsMap,
targetAssignment, preChecksResult, summaryResult);
}
}
}
Expand Down
Loading