Skip to content

Commit 472195e

Browse files
pthirunclaude
andcommitted
[controller][common][protocol] Add store-level ULE config for RT topics
Add uncleanLeaderElectionEnabledForRTTopics as a store-level config using the tri-state pattern (NOT_SPECIFIED/ENABLED/DISABLED). When NOT_SPECIFIED, falls back to the cluster-level config. This enables per-store tracking of ULE settings and preserves the setting during store migration. Changes: - New Avro schema versions (StoreMetaValue v41, AdminOperation v96) - Store interface/impl (Store, ZKStore, ReadOnlyStore, SystemStore, StoreInfo) - UpdateStoreQueryParams with migration constructor support - Controller logic (VeniceParentHelixAdmin, AdminExecutionTask, VeniceHelixAdmin) - RealTimeTopicSwitcher store-level override with cluster-level fallback - resolveUncleanLeaderElection helper for store-then-cluster resolution - Tests for override, fallback, and resolution logic Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 2f0d39c commit 472195e

19 files changed

Lines changed: 2003 additions & 7 deletions

File tree

build.gradle

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -327,8 +327,8 @@ subprojects {
327327
// when actually using the new protocol. Example to pin KME to v12 when introducing v13:
328328
// project(':internal:venice-common').file('src/main/resources/avro/KafkaMessageEnvelope/v12', PathValidation.DIRECTORY)
329329
def versionOverrides = [
330-
project(':internal:venice-common').file('src/main/resources/avro/StoreMetaValue/v39', PathValidation.DIRECTORY),
331-
project(':services:venice-controller').file('src/main/resources/avro/AdminOperation/v94', PathValidation.DIRECTORY)
330+
project(':internal:venice-common').file('src/main/resources/avro/StoreMetaValue/v41', PathValidation.DIRECTORY),
331+
project(':services:venice-controller').file('src/main/resources/avro/AdminOperation/v96', PathValidation.DIRECTORY)
332332
]
333333

334334
def schemaDirs = [sourceDir]

internal/venice-common/src/main/java/com/linkedin/venice/controllerapi/ControllerApiConstants.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -268,6 +268,8 @@ public class ControllerApiConstants {
268268

269269
public static final String BLOB_TRANSFER_ENABLED = "blob_transfer_enabled";
270270
public static final String BLOB_TRANSFER_IN_SERVER_ENABLED = "blob_transfer_in_server_enabled";
271+
public static final String UNCLEAN_LEADER_ELECTION_ENABLED_FOR_RT_TOPICS =
272+
"unclean_leader_election_enabled_for_rt_topics";
271273

272274
public static final String HEARTBEAT_TIMESTAMP = "heartbeat_timestamp";
273275

internal/venice-common/src/main/java/com/linkedin/venice/controllerapi/UpdateStoreQueryParams.java

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -78,6 +78,7 @@
7878
import static com.linkedin.venice.controllerapi.ControllerApiConstants.TARGET_SWAP_REGION_WAIT_TIME;
7979
import static com.linkedin.venice.controllerapi.ControllerApiConstants.TIME_LAG_TO_GO_ONLINE;
8080
import static com.linkedin.venice.controllerapi.ControllerApiConstants.TTL_REPUSH_ENABLED;
81+
import static com.linkedin.venice.controllerapi.ControllerApiConstants.UNCLEAN_LEADER_ELECTION_ENABLED_FOR_RT_TOPICS;
8182
import static com.linkedin.venice.controllerapi.ControllerApiConstants.UNUSED_SCHEMA_DELETION_ENABLED;
8283
import static com.linkedin.venice.controllerapi.ControllerApiConstants.UPDATED_CONFIGS_LIST;
8384
import static com.linkedin.venice.controllerapi.ControllerApiConstants.VERSION;
@@ -162,6 +163,8 @@ public UpdateStoreQueryParams(StoreInfo srcStore, boolean storeMigrating) {
162163
.setBlobTransferEnabled(srcStore.isBlobTransferEnabled())
163164
.setBlobTransferInServerEnabled(
164165
ConfigCommonUtils.ActivationState.valueOf(srcStore.getBlobTransferInServerEnabled()))
166+
.setUncleanLeaderElectionEnabledForRTTopics(
167+
ConfigCommonUtils.ActivationState.valueOf(srcStore.getUncleanLeaderElectionEnabledForRTTopics()))
165168
.setMaxRecordSizeBytes(srcStore.getMaxRecordSizeBytes())
166169
.setMaxNearlineRecordSizeBytes(srcStore.getMaxNearlineRecordSizeBytes())
167170
.setTargetRegionSwap(srcStore.getTargetRegionSwap())
@@ -814,6 +817,14 @@ public Optional<String> getBlobTransferInServerEnabled() {
814817
return getString(BLOB_TRANSFER_IN_SERVER_ENABLED);
815818
}
816819

820+
public UpdateStoreQueryParams setUncleanLeaderElectionEnabledForRTTopics(ConfigCommonUtils.ActivationState state) {
821+
return putString(UNCLEAN_LEADER_ELECTION_ENABLED_FOR_RT_TOPICS, state.name());
822+
}
823+
824+
public Optional<String> getUncleanLeaderElectionEnabledForRTTopics() {
825+
return getString(UNCLEAN_LEADER_ELECTION_ENABLED_FOR_RT_TOPICS);
826+
}
827+
817828
public UpdateStoreQueryParams setNearlineProducerCompressionEnabled(boolean compressionEnabled) {
818829
return putBoolean(NEARLINE_PRODUCER_COMPRESSION_ENABLED, compressionEnabled);
819830
}

internal/venice-common/src/main/java/com/linkedin/venice/meta/ReadOnlyStore.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1650,6 +1650,16 @@ public String getBlobTransferInServerEnabled() {
16501650
return this.delegate.getBlobTransferInServerEnabled();
16511651
}
16521652

1653+
@Override
1654+
public String getUncleanLeaderElectionEnabledForRTTopics() {
1655+
return this.delegate.getUncleanLeaderElectionEnabledForRTTopics();
1656+
}
1657+
1658+
@Override
1659+
public void setUncleanLeaderElectionEnabledForRTTopics(String uncleanLeaderElectionEnabledForRTTopics) {
1660+
throw new UnsupportedOperationException("Unclean leader election config is read-only");
1661+
}
1662+
16531663
@Override
16541664
public boolean isNearlineProducerCompressionEnabled() {
16551665
return delegate.isNearlineProducerCompressionEnabled();

internal/venice-common/src/main/java/com/linkedin/venice/meta/Store.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -340,6 +340,10 @@ default IntSet getVersionNumbers() {
340340

341341
String getBlobTransferInServerEnabled();
342342

343+
String getUncleanLeaderElectionEnabledForRTTopics();
344+
345+
void setUncleanLeaderElectionEnabledForRTTopics(String uncleanLeaderElectionEnabledForRTTopics);
346+
343347
boolean isNearlineProducerCompressionEnabled();
344348

345349
void setNearlineProducerCompressionEnabled(boolean compressionEnabled);

internal/venice-common/src/main/java/com/linkedin/venice/meta/StoreInfo.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -81,6 +81,7 @@ public static StoreInfo fromStore(Store store) {
8181
storeInfo.setUnusedSchemaDeletionEnabled(store.isUnusedSchemaDeletionEnabled());
8282
storeInfo.setBlobTransferEnabled(store.isBlobTransferEnabled());
8383
storeInfo.setBlobTransferInServerEnabled(store.getBlobTransferInServerEnabled());
84+
storeInfo.setUncleanLeaderElectionEnabledForRTTopics(store.getUncleanLeaderElectionEnabledForRTTopics());
8485
storeInfo.setNearlineProducerCompressionEnabled(store.isNearlineProducerCompressionEnabled());
8586
storeInfo.setNearlineProducerCountPerWriter(store.getNearlineProducerCountPerWriter());
8687
storeInfo.setTargetRegionSwap(store.getTargetSwapRegion());
@@ -361,6 +362,7 @@ public static StoreInfo fromStore(Store store) {
361362

362363
private boolean blobTransferEnabled;
363364
private String blobTransferInServerEnable = ActivationState.NOT_SPECIFIED.name();
365+
private String uncleanLeaderElectionEnabledForRTTopics = ActivationState.NOT_SPECIFIED.name();
364366

365367
private boolean nearlineProducerCompressionEnabled;
366368
private int nearlineProducerCountPerWriter;
@@ -900,6 +902,14 @@ public String getBlobTransferInServerEnabled() {
900902
return this.blobTransferInServerEnable;
901903
}
902904

905+
public void setUncleanLeaderElectionEnabledForRTTopics(String uncleanLeaderElectionEnabledForRTTopics) {
906+
this.uncleanLeaderElectionEnabledForRTTopics = uncleanLeaderElectionEnabledForRTTopics;
907+
}
908+
909+
public String getUncleanLeaderElectionEnabledForRTTopics() {
910+
return this.uncleanLeaderElectionEnabledForRTTopics;
911+
}
912+
903913
public boolean isNearlineProducerCompressionEnabled() {
904914
return nearlineProducerCompressionEnabled;
905915
}

internal/venice-common/src/main/java/com/linkedin/venice/meta/SystemStore.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -707,6 +707,16 @@ public String getBlobTransferInServerEnabled() {
707707
return zkSharedStore.getBlobTransferInServerEnabled();
708708
}
709709

710+
@Override
711+
public String getUncleanLeaderElectionEnabledForRTTopics() {
712+
return zkSharedStore.getUncleanLeaderElectionEnabledForRTTopics();
713+
}
714+
715+
@Override
716+
public void setUncleanLeaderElectionEnabledForRTTopics(String uncleanLeaderElectionEnabledForRTTopics) {
717+
throwUnsupportedOperationException("setUncleanLeaderElectionEnabledForRTTopics is not supported in SystemStore");
718+
}
719+
710720
@Override
711721
public void setMaxCompactionLagSeconds(long maxCompactionLagSeconds) {
712722
throwUnsupportedOperationException("setMaxCompactionLagSeconds");

internal/venice-common/src/main/java/com/linkedin/venice/meta/ZKStore.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -976,6 +976,16 @@ public String getBlobTransferInServerEnabled() {
976976
return this.storeProperties.blobTransferInServerEnabled.toString();
977977
}
978978

979+
@Override
980+
public String getUncleanLeaderElectionEnabledForRTTopics() {
981+
return this.storeProperties.uncleanLeaderElectionEnabledForRTTopics.toString();
982+
}
983+
984+
@Override
985+
public void setUncleanLeaderElectionEnabledForRTTopics(String uncleanLeaderElectionEnabledForRTTopics) {
986+
this.storeProperties.uncleanLeaderElectionEnabledForRTTopics = uncleanLeaderElectionEnabledForRTTopics;
987+
}
988+
979989
@Override
980990
public boolean isNearlineProducerCompressionEnabled() {
981991
return this.storeProperties.nearlineProducerCompressionEnabled;

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,7 @@ public enum AvroProtocolDefinition {
7777
*
7878
* TODO: Move AdminOperation to venice-common module so that we can properly reference it here.
7979
*/
80-
ADMIN_OPERATION(94, SpecificData.get().getSchema(ByteBuffer.class), "AdminOperation"),
80+
ADMIN_OPERATION(96, SpecificData.get().getSchema(ByteBuffer.class), "AdminOperation"),
8181

8282
/**
8383
* Single chunk of a large multi-chunk value. Just a bunch of bytes.
@@ -148,7 +148,7 @@ public enum AvroProtocolDefinition {
148148
/**
149149
* Value schema for metadata system store.
150150
*/
151-
METADATA_SYSTEM_SCHEMA_STORE(39, StoreMetaValue.class),
151+
METADATA_SYSTEM_SCHEMA_STORE(41, StoreMetaValue.class),
152152

153153
/*
154154
Value Schema for Parent Controller Metadata system store

0 commit comments

Comments
 (0)