Skip to content

Commit b810be5

Browse files
KaiSernLimCopilot
andcommitted
Remove MARK_VERSION_ROLLED_BACK admin operation
Replace cross-region protocol message with direct deleteOldVersion calls per region via controllerClientMap. Non-current PUSHED/ONLINE child copies of a terminal deferred version are deleted immediately rather than transitioned through ROLLED_BACK; the double-guard in VeniceHelixAdmin and controllerClientMap ensures current versions are never touched. Co-authored-by: Copilot App <223556219+Copilot@users.noreply.github.com>
1 parent 6de7696 commit b810be5

12 files changed

Lines changed: 73 additions & 1700 deletions

File tree

docs/operations/data-management/version-lifecycle.md

Lines changed: 11 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -8,9 +8,9 @@ backup, within configured retention, or part of a store migration.
88
## Deletion decision table
99

1010
The count-based sweep is implemented by `Store.retrieveVersionsToDelete`. The time-based sweep is
11-
implemented by `StoreBackupVersionCleanupService`. Deferred version swap terminal transitions first
12-
reconcile bootstrap-complete, non-current child copies to `ROLLED_BACK`, so they cannot remain
13-
invisible to both sweeps.
11+
implemented by `StoreBackupVersionCleanupService`. Deferred version swap terminal transitions
12+
directly delete bootstrap-complete, non-current child copies via `ControllerClient.deleteOldVersion`,
13+
so they cannot remain invisible to both sweeps.
1414

1515
| Initial status | Role / initial condition | Count retention | Time / safety gate | Trigger | Deletion decision |
1616
| --- | --- | --- | --- | --- | --- |
@@ -26,7 +26,7 @@ invisible to both sweeps.
2626
| `PUSHED` | Non-current version without a terminal parent decision | Excluded because it may be an active deferred or concurrent swap candidate | Not eligible | Count or time cleanup | **KEEP** |
2727
| `PUSHED` or `ONLINE` | Still current in a child after parent `ERROR` or `ROLLED_BACK` | Protected while serving; reconciliation remains incomplete | Not eligible | Parent terminal reconciliation | **DEFER: keep and retry** |
2828
| `PUSHED` or `ONLINE` | Still current in a child after parent `PARTIALLY_ONLINE` | Serving the intentional partial state | Not eligible | Parent terminal reconciliation | **KEEP** |
29-
| `PUSHED` or non-current `ONLINE` | Parent deferred swap becomes `ERROR`, `PARTIALLY_ONLINE`, or `ROLLED_BACK` | Removed from count sweep | Rolled-back retention gate applies | Parent terminal transition | **DEFER: mark `ROLLED_BACK`** |
29+
| `PUSHED` or non-current `ONLINE` | Parent deferred swap becomes `ERROR`, `PARTIALLY_ONLINE`, or `ROLLED_BACK` | Removed from count sweep | None — immediately deleted | Parent terminal transition | **DELETE** |
3030
| `ONLINE` | Backup within the configured preserved count | Within limit | Not considered | Count sweep | **KEEP** |
3131
| `ONLINE` | Backup beyond the configured preserved count | Exceeds limit | Not considered | Count sweep | **DELETE** |
3232
| `ONLINE` | Backup considered by retention cleanup | Not considered | Before minimum retention | Time sweep | **DEFER** |
@@ -39,15 +39,11 @@ invisible to both sweeps.
3939
Count-based and time-based cleanup are independent triggers. The first applicable trigger may delete
4040
an eligible backup, but neither trigger may delete the current version. `PUSHED` versions are never
4141
deleted from status and version number alone because multiple deferred or concurrent swaps can be
42-
active. A terminal parent decision explicitly converts bootstrap-complete non-current child copies
43-
to `ROLLED_BACK`. The controller scans every deferred terminal parent version, including versions
44-
superseded by a newer push, and retries while a child is unreachable, missing the target metadata,
45-
in progress, or unexpectedly still current after parent `ERROR` or `ROLLED_BACK`.
46-
47-
When a child copy becomes `ROLLED_BACK`, the controller resets the store-level latest-promotion
48-
timestamp to start the rollback retention window. A subsequent promotion can reset that shared clock
49-
again, so a rolled-back version may be retained longer than the configured duration; it cannot be
50-
deleted immediately because the current version was promoted long before the rollback.
42+
active. A terminal parent decision explicitly deletes bootstrap-complete non-current child copies
43+
via per-region `ControllerClient.deleteOldVersion`. The controller scans every deferred terminal
44+
parent version, including versions superseded by a newer push, and retries while a child is
45+
unreachable, missing the target metadata, in progress, or unexpectedly still current after parent
46+
`ERROR` or `ROLLED_BACK`.
5147

5248
## State machine
5349

@@ -60,7 +56,7 @@ stateDiagram-v2
6056
STARTED --> ERROR: push fails
6157
STARTED --> KILLED: push is killed
6258
PUSHED --> ONLINE: version swap succeeds
63-
PUSHED --> ROLLED_BACK: deferred swap terminates while non-current
59+
PUSHED --> DELETED: deferred swap terminates while non-current (parent terminal sweep)
6460
ONLINE --> ROLLED_BACK: rollback or abandoned non-current copy
6561
6662
NOT_CREATED --> DELETED: stale below-current metadata after time gate
@@ -81,7 +77,7 @@ stateDiagram-v2
8177
Before terminal parent decision: KEEP
8278
Current after ERROR/ROLLED_BACK: retry
8379
Current after PARTIALLY_ONLINE: KEEP
84-
Parent terminal + non-current: ROLLED_BACK
80+
Parent terminal + non-current: DELETE immediately
8581
end note
8682
8783
note right of ONLINE

internal/venice-common/src/main/java/com/linkedin/venice/serialization/avro/AvroProtocolDefinition.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,7 @@ public enum AvroProtocolDefinition {
7171
*
7272
* TODO: Move AdminOperation to venice-common module so that we can properly reference it here.
7373
*/
74-
ADMIN_OPERATION(102, SpecificData.get().getSchema(ByteBuffer.class), "AdminOperation"),
74+
ADMIN_OPERATION(101, SpecificData.get().getSchema(ByteBuffer.class), "AdminOperation"),
7575

7676
/**
7777
* Single chunk of a large multi-chunk value. Just a bunch of bytes.

services/venice-controller/src/main/java/com/linkedin/venice/controller/DeferredVersionSwapService.java

Lines changed: 20 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1418,14 +1418,26 @@ private boolean reconcileAbandonedVersionInChildRegions(
14181418
}
14191419

14201420
if (!regionsToReconcile.isEmpty()) {
1421-
String regionsFilter = RegionUtils.composeRegionList(regionsToReconcile);
1422-
veniceParentHelixAdmin.markVersionRolledBack(clusterName, storeName, targetVersionNum, regionsFilter);
1423-
LOGGER.info(
1424-
"Reconciled version: {} for store: {} as ROLLED_BACK in child regions: {} after parent status: {}",
1425-
targetVersionNum,
1426-
storeName,
1427-
regionsFilter,
1428-
parentStatus);
1421+
for (String region: regionsToReconcile) {
1422+
try {
1423+
controllerClientMap.get(region).deleteOldVersion(storeName, targetVersionNum);
1424+
LOGGER.info(
1425+
"Deleted abandoned deferred version {} for store {} in region {} (parent status: {})",
1426+
targetVersionNum,
1427+
storeName,
1428+
region,
1429+
parentStatus);
1430+
} catch (Exception e) {
1431+
LOGGER.warn(
1432+
"Failed to delete abandoned deferred version {} for store {} in region {} (parent status: {})",
1433+
targetVersionNum,
1434+
storeName,
1435+
region,
1436+
parentStatus,
1437+
e);
1438+
allRegionsTerminal = false;
1439+
}
1440+
}
14291441
}
14301442
return allRegionsTerminal;
14311443
}

services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceHelixAdmin.java

Lines changed: 0 additions & 70 deletions
Original file line numberDiff line numberDiff line change
@@ -5124,76 +5124,6 @@ public void rollbackToBackupVersion(String clusterName, String storeName, String
51245124
}
51255125
}
51265126

5127-
/**
5128-
* Mark {@code versionNumber} as {@link VersionStatus#ROLLED_BACK} in the regions selected by
5129-
* {@code regionFilter}, WITHOUT changing the store's current version.
5130-
*
5131-
* <p>Deferred swaps can end globally while a bootstrap-complete copy is still non-current in some
5132-
* regions. Marking that copy ROLLED_BACK makes it eligible for time-based cleanup. In-progress
5133-
* versions and the current serving version remain untouched.
5134-
*/
5135-
public void markVersionRolledBack(String clusterName, String storeName, int versionNumber, String regionFilter) {
5136-
String currentRegion = getRegionName();
5137-
if (!StringUtils.isEmpty(regionFilter)) {
5138-
Set<String> regionsFilter = parseRegionsFilterList(regionFilter);
5139-
if (!regionsFilter.contains(currentRegion)) {
5140-
LOGGER.info(
5141-
"markVersionRolledBack will be skipped for store: {} version: {} in cluster: {}, because the region filter"
5142-
+ " is {} which doesn't include the current region: {}",
5143-
storeName,
5144-
versionNumber,
5145-
clusterName,
5146-
regionsFilter,
5147-
currentRegion);
5148-
return;
5149-
}
5150-
}
5151-
5152-
storeMetadataUpdate(clusterName, storeName, (store, resources) -> {
5153-
if (!store.containsVersion(versionNumber)) {
5154-
LOGGER.info(
5155-
"markVersionRolledBack skipped: version {} not found in store {} in cluster {}",
5156-
versionNumber,
5157-
storeName,
5158-
clusterName);
5159-
return store;
5160-
}
5161-
// Never override the current version's status; that version serves reads.
5162-
if (store.getCurrentVersion() == versionNumber) {
5163-
LOGGER.warn(
5164-
"markVersionRolledBack skipped: version {} is the current version of store {} in cluster {}",
5165-
versionNumber,
5166-
storeName,
5167-
clusterName);
5168-
return store;
5169-
}
5170-
VersionStatus currentStatus = store.getVersionStatus(versionNumber);
5171-
if (VersionStatus.isVersionRolledBack(currentStatus) || VersionStatus.canDelete(currentStatus)) {
5172-
// Already ROLLED_BACK, or already in a terminal cleanable status (KILLED/ERROR) — nothing to do.
5173-
return store;
5174-
}
5175-
if (!VersionStatus.isBootstrapCompleted(currentStatus)) {
5176-
LOGGER.info(
5177-
"markVersionRolledBack skipped: version {} of store {} in cluster {} has not completed bootstrap (status {})",
5178-
versionNumber,
5179-
storeName,
5180-
clusterName,
5181-
currentStatus);
5182-
return store;
5183-
}
5184-
LOGGER.info(
5185-
"Marking version {} of store {} in cluster {} as ROLLED_BACK (was {})",
5186-
versionNumber,
5187-
storeName,
5188-
clusterName,
5189-
currentStatus);
5190-
store.updateVersionStatus(versionNumber, ROLLED_BACK);
5191-
store.setLatestVersionPromoteToCurrentTimestamp(
5192-
Math.max(store.getLatestVersionPromoteToCurrentTimestamp(), System.currentTimeMillis()));
5193-
return store;
5194-
});
5195-
}
5196-
51975127
/**
51985128
* Update the largest used version number of a specified store.
51995129
*/

services/venice-controller/src/main/java/com/linkedin/venice/controller/VeniceParentHelixAdmin.java

Lines changed: 0 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,6 @@
6363
import com.linkedin.venice.controller.kafka.protocol.admin.PushStatusSystemStoreAutoCreationValidation;
6464
import com.linkedin.venice.controller.kafka.protocol.admin.ResumeStore;
6565
import com.linkedin.venice.controller.kafka.protocol.admin.RollbackCurrentVersion;
66-
import com.linkedin.venice.controller.kafka.protocol.admin.MarkVersionRolledBack;
6766
import com.linkedin.venice.controller.kafka.protocol.admin.SchemaMeta;
6867
import com.linkedin.venice.controller.kafka.protocol.admin.SetStoreOwner;
6968
import com.linkedin.venice.controller.kafka.protocol.admin.SetStorePartitionCount;
@@ -2520,34 +2519,6 @@ void checkNewPushCapacityFromChildren(String clusterName, String storeName) {
25202519
}
25212520
}
25222521

2523-
/**
2524-
* Mark {@code versionNum} as ROLLED_BACK in the child regions selected by {@code regionFilter},
2525-
* WITHOUT changing the current version. Used to reconcile the non-target regions after a deferred
2526-
* version swap rollback: those regions bootstrapped the target version but never swapped to it, so
2527-
* their copy is stranded in PUSHED status above their unchanged current version and leaks disk.
2528-
* Marking it ROLLED_BACK makes it visible to {@code StoreBackupVersionCleanupService}.
2529-
*/
2530-
public void markVersionRolledBack(String clusterName, String storeName, int versionNum, String regionFilter) {
2531-
acquireAdminMessageLock(clusterName, storeName);
2532-
try {
2533-
getVeniceHelixAdmin().checkPreConditionForUpdateStoreMetadata(clusterName, storeName);
2534-
2535-
MarkVersionRolledBack markVersionRolledBack =
2536-
(MarkVersionRolledBack) AdminMessageType.MARK_VERSION_ROLLED_BACK.getNewInstance();
2537-
markVersionRolledBack.clusterName = clusterName;
2538-
markVersionRolledBack.storeName = storeName;
2539-
markVersionRolledBack.versionNum = versionNum;
2540-
markVersionRolledBack.regionsFilter = regionFilter;
2541-
AdminOperation message = new AdminOperation();
2542-
message.operationType = AdminMessageType.MARK_VERSION_ROLLED_BACK.getValue();
2543-
message.payloadUnion = markVersionRolledBack;
2544-
2545-
sendAdminMessageAndWaitForConsumed(clusterName, storeName, message);
2546-
} finally {
2547-
releaseAdminMessageLock(clusterName, storeName);
2548-
}
2549-
}
2550-
25512522
private void updateParentVersionStatusAfterRollback(
25522523
String clusterName,
25532524
String storeName,

services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/consumer/AdminExecutionTask.java

Lines changed: 0 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,6 @@
3030
import com.linkedin.venice.controller.kafka.protocol.admin.ResumeStore;
3131
import com.linkedin.venice.controller.kafka.protocol.admin.RollForwardCurrentVersion;
3232
import com.linkedin.venice.controller.kafka.protocol.admin.RollbackCurrentVersion;
33-
import com.linkedin.venice.controller.kafka.protocol.admin.MarkVersionRolledBack;
3433
import com.linkedin.venice.controller.kafka.protocol.admin.SetStoreCurrentVersion;
3534
import com.linkedin.venice.controller.kafka.protocol.admin.SetStoreOwner;
3635
import com.linkedin.venice.controller.kafka.protocol.admin.SetStorePartitionCount;
@@ -345,9 +344,6 @@ private void processMessage(AdminOperation adminOperation) {
345344
case ROLLFORWARD_CURRENT_VERSION:
346345
handleRollForwardToFutureVersion((RollForwardCurrentVersion) adminOperation.payloadUnion);
347346
break;
348-
case MARK_VERSION_ROLLED_BACK:
349-
handleMarkVersionRolledBack((MarkVersionRolledBack) adminOperation.payloadUnion);
350-
break;
351347
default:
352348
throw new VeniceException("Unknown admin operation type: " + adminOperation.operationType);
353349
}
@@ -939,15 +935,6 @@ private void handleRollbackCurrentVersion(RollbackCurrentVersion message) {
939935
admin.rollbackToBackupVersion(clusterName, storeName, regionFilter);
940936
}
941937

942-
private void handleMarkVersionRolledBack(MarkVersionRolledBack message) {
943-
String clusterName = message.getClusterName().toString();
944-
String storeName = message.getStoreName().toString();
945-
int versionNum = message.getVersionNum();
946-
CharSequence regionsFilter = message.getRegionsFilter();
947-
String regionFilter = regionsFilter == null ? null : regionsFilter.toString();
948-
admin.markVersionRolledBack(clusterName, storeName, versionNum, regionFilter);
949-
}
950-
951938
private void handleDeleteUnusedValueSchema(DeleteUnusedValueSchemas message) {
952939
String clusterName = message.getClusterName().toString();
953940
String storeName = message.getStoreName().toString();

services/venice-controller/src/main/java/com/linkedin/venice/controller/kafka/protocol/enums/AdminMessageType.java

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,6 @@
1515
import com.linkedin.venice.controller.kafka.protocol.admin.DisableStoreRead;
1616
import com.linkedin.venice.controller.kafka.protocol.admin.EnableStoreRead;
1717
import com.linkedin.venice.controller.kafka.protocol.admin.KillOfflinePushJob;
18-
import com.linkedin.venice.controller.kafka.protocol.admin.MarkVersionRolledBack;
1918
import com.linkedin.venice.controller.kafka.protocol.admin.MetaSystemStoreAutoCreationValidation;
2019
import com.linkedin.venice.controller.kafka.protocol.admin.MetadataSchemaCreation;
2120
import com.linkedin.venice.controller.kafka.protocol.admin.MigrateStore;
@@ -66,8 +65,7 @@ public enum AdminMessageType implements VeniceDimensionInterface {
6665
CONFIGURE_INCREMENTAL_PUSH_FOR_CLUSTER(22, true), META_SYSTEM_STORE_AUTO_CREATION_VALIDATION(23, false),
6766
PUSH_STATUS_SYSTEM_STORE_AUTO_CREATION_VALIDATION(24, false), CREATE_STORAGE_PERSONA(25, false),
6867
DELETE_STORAGE_PERSONA(26, false), UPDATE_STORAGE_PERSONA(27, false), DELETE_UNUSED_VALUE_SCHEMA(28, false),
69-
ROLLBACK_CURRENT_VERSION(29, false), ROLLFORWARD_CURRENT_VERSION(30, false),
70-
MARK_VERSION_ROLLED_BACK(31, false);
68+
ROLLBACK_CURRENT_VERSION(29, false), ROLLFORWARD_CURRENT_VERSION(30, false);
7169

7270
private final int value;
7371
private final boolean batchUpdate;
@@ -140,8 +138,6 @@ public Object getNewInstance() {
140138
return new RollForwardCurrentVersion();
141139
case ROLLBACK_CURRENT_VERSION:
142140
return new RollbackCurrentVersion();
143-
case MARK_VERSION_ROLLED_BACK:
144-
return new MarkVersionRolledBack();
145141
default:
146142
throw new VeniceException("Unsupported " + getClass().getSimpleName() + " value: " + value);
147143
}

0 commit comments

Comments
 (0)