Skip to content

Commit d91fdce

Browse files
Sanchit JainCopilot
andcommitted
Rewire minimum active replicas to IdealState and remove the ResourceConfig copy
Remove only MIN_ACTIVE_REPLICAS from ResourceConfig and read the minimum from IdealState in delayed and WAGED rebalance paths, retaining the replica-count fallback. Preserve all other ResourceConfig settings. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
1 parent 7487435 commit d91fdce

7 files changed

Lines changed: 19 additions & 66 deletions

File tree

helix-core/src/main/java/org/apache/helix/controller/changedetector/trimmer/ResourceConfigTrimmer.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ public class ResourceConfigTrimmer extends HelixPropertyTrimmer<ResourceConfig>
4040
* MONITORING_DISABLED,
4141
* STATE_MODEL_DEF_REF,
4242
* REPLICAS,
43+
* MIN_ACTIVE_REPLICAS,
4344
* EXTERNAL_VIEW_DISABLED,
4445
* DELAY_REBALANCE_ENABLED,
4546
* HELIX_ENABLED
@@ -48,7 +49,6 @@ public class ResourceConfigTrimmer extends HelixPropertyTrimmer<ResourceConfig>
4849
.of(FieldType.SIMPLE_FIELD, ImmutableSet
4950
.of(ResourceConfigProperty.NUM_PARTITIONS.name(),
5051
ResourceConfigProperty.STATE_MODEL_FACTORY_NAME.name(),
51-
ResourceConfigProperty.MIN_ACTIVE_REPLICAS.name(),
5252
ResourceConfigProperty.MAX_PARTITIONS_PER_INSTANCE.name(),
5353
ResourceConfigProperty.INSTANCE_GROUP_TAG.name()),
5454
FieldType.MAP_FIELD, ImmutableSet

helix-core/src/main/java/org/apache/helix/controller/rebalancer/DelayedAutoRebalancer.java

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -171,9 +171,8 @@ public IdealState computeNewIdealState(String resourceName,
171171
|| liveEnabledAssignableNodeList.size() != activeNodes.size()) {
172172
List<String> activeNodeList = new ArrayList<>(activeNodes);
173173
Collections.sort(activeNodeList);
174-
int minActiveReplicas = DelayedRebalanceUtil.getMinActiveReplica(
175-
ResourceConfig.mergeIdealStateWithResourceConfig(resourceConfig, currentIdealState),
176-
currentIdealState, replicaCount);
174+
int minActiveReplicas =
175+
DelayedRebalanceUtil.getMinActiveReplica(currentIdealState, replicaCount);
177176

178177
ZNRecord newActiveMapping =
179178
_rebalanceStrategy.computePartitionAssignment(allNodeList, activeNodeList, currentMapping,

helix-core/src/main/java/org/apache/helix/controller/rebalancer/util/DelayedRebalanceUtil.java

Lines changed: 8 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -235,25 +235,17 @@ public static Map<String, List<String>> getFinalDelayedMapping(
235235
/**
236236
* Get the minimum active replica count threshold that allows delayed rebalance.
237237
* Prioritize of the input params:
238-
* 1. resourceConfig
239-
* 2. idealState
240-
* 3. replicaCount
238+
* 1. idealState
239+
* 2. replicaCount
241240
* The lower priority minimum active replica count will only be applied if the higher priority
242-
* items are missing.
243-
* TODO: Remove the idealState input once we have all the config information migrated to the
244-
* TODO: resource config by default.
241+
* item is missing.
245242
*
246-
* @param resourceConfig the resource config
247243
* @param idealState the ideal state of the resource
248244
* @param replicaCount the expected active replica count.
249245
* @return the expected minimum active replica count that is required
250246
*/
251-
public static int getMinActiveReplica(ResourceConfig resourceConfig, IdealState idealState,
252-
int replicaCount) {
253-
int minActiveReplicas = resourceConfig == null ? -1 : resourceConfig.getMinActiveReplica();
254-
if (minActiveReplicas < 0) {
255-
minActiveReplicas = idealState.getMinActiveReplicas();
256-
}
247+
public static int getMinActiveReplica(IdealState idealState, int replicaCount) {
248+
int minActiveReplicas = idealState.getMinActiveReplicas();
257249
if (minActiveReplicas < 0) {
258250
minActiveReplicas = replicaCount;
259251
}
@@ -408,9 +400,8 @@ private static List<String> findPartitionsMissingMinActiveReplica(
408400
IdealState currentIdealState = clusterData.getIdealState(resourceName);
409401
Set<String> enabledLiveInstances = clusterData.getEnabledLiveInstances();
410402
int numReplica = currentIdealState.getReplicaCount(enabledLiveInstances.size());
411-
int minActiveReplica = DelayedRebalanceUtil.getMinActiveReplica(ResourceConfig
412-
.mergeIdealStateWithResourceConfig(clusterData.getResourceConfig(resourceName),
413-
currentIdealState), currentIdealState, numReplica);
403+
int minActiveReplica =
404+
DelayedRebalanceUtil.getMinActiveReplica(currentIdealState, numReplica);
414405
return resourceAssignment.getMappedPartitions()
415406
.parallelStream()
416407
.filter(partition -> {
@@ -429,9 +420,7 @@ private static int getMinActiveReplica(ResourceControllerDataProvider clusterDat
429420
IdealState currentIdealState = clusterData.getIdealState(resourceName);
430421
Set<String> enabledLiveInstances = clusterData.getEnabledLiveInstances();
431422
int numReplica = currentIdealState.getReplicaCount(enabledLiveInstances.size());
432-
return DelayedRebalanceUtil.getMinActiveReplica(ResourceConfig
433-
.mergeIdealStateWithResourceConfig(clusterData.getResourceConfig(resourceName),
434-
currentIdealState), currentIdealState, numReplica);
423+
return DelayedRebalanceUtil.getMinActiveReplica(currentIdealState, numReplica);
435424
}
436425

437426
/**

helix-core/src/main/java/org/apache/helix/controller/rebalancer/waged/WagedRebalancer.java

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -898,9 +898,8 @@ protected boolean requireRebalanceOverwrite(ResourceControllerDataProvider clust
898898
Set<String> enabledLiveInstances = clusterData.getEnabledLiveInstances();
899899

900900
int numReplica = currentIdealState.getReplicaCount(enabledLiveInstances.size());
901-
int minActiveReplica = DelayedRebalanceUtil.getMinActiveReplica(ResourceConfig
902-
.mergeIdealStateWithResourceConfig(clusterData.getResourceConfig(resourceName),
903-
currentIdealState), currentIdealState, numReplica);
901+
int minActiveReplica =
902+
DelayedRebalanceUtil.getMinActiveReplica(currentIdealState, numReplica);
904903
resourceAssignment.getMappedPartitions().parallelStream().forEach(partition -> {
905904
int enabledLivePlacementCounter = 0;
906905
for (String instance : resourceAssignment.getReplicaMap(partition).keySet()) {

helix-core/src/main/java/org/apache/helix/model/ResourceConfig.java

Lines changed: 4 additions & 31 deletions
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,6 @@ public enum ResourceConfigProperty {
4949
MONITORING_DISABLED, // Resource-level config, do not create Mbean and report any status for the resource.
5050
NUM_PARTITIONS,
5151
STATE_MODEL_FACTORY_NAME,
52-
MIN_ACTIVE_REPLICAS,
5352
MAX_PARTITIONS_PER_INSTANCE,
5453
INSTANCE_GROUP_TAG,
5554
HELIX_ENABLED,
@@ -99,20 +98,20 @@ public ResourceConfig(ZNRecord record, String id) {
9998

10099
public ResourceConfig(String resourceId, Boolean monitorDisabled, int numPartitions,
101100
String stateModelFactoryName,
102-
int minActiveReplica, int maxPartitionsPerInstance, String instanceGroupTag,
101+
int maxPartitionsPerInstance, String instanceGroupTag,
103102
Boolean helixEnabled, Boolean externalViewDisabled, RebalanceConfig rebalanceConfig,
104103
StateTransitionTimeoutConfig stateTransitionTimeoutConfig,
105104
Map<String, List<String>> listFields, Map<String, Map<String, String>> mapFields,
106105
Boolean p2pMessageEnabled) {
107106
this(resourceId, monitorDisabled, numPartitions, stateModelFactoryName,
108-
minActiveReplica, maxPartitionsPerInstance, instanceGroupTag, helixEnabled,
107+
maxPartitionsPerInstance, instanceGroupTag, helixEnabled,
109108
externalViewDisabled, rebalanceConfig, stateTransitionTimeoutConfig, listFields, mapFields,
110109
p2pMessageEnabled, null);
111110
}
112111

113112
private ResourceConfig(String resourceId, Boolean monitorDisabled, int numPartitions,
114113
String stateModelFactoryName,
115-
int minActiveReplica, int maxPartitionsPerInstance, String instanceGroupTag,
114+
int maxPartitionsPerInstance, String instanceGroupTag,
116115
Boolean helixEnabled, Boolean externalViewDisabled, RebalanceConfig rebalanceConfig,
117116
StateTransitionTimeoutConfig stateTransitionTimeoutConfig,
118117
Map<String, List<String>> listFields, Map<String, Map<String, String>> mapFields,
@@ -135,10 +134,6 @@ private ResourceConfig(String resourceId, Boolean monitorDisabled, int numPartit
135134
_record.setSimpleField(ResourceConfigProperty.STATE_MODEL_FACTORY_NAME.name(), stateModelFactoryName);
136135
}
137136

138-
if (minActiveReplica >= 0) {
139-
_record.setIntField(ResourceConfigProperty.MIN_ACTIVE_REPLICAS.name(), minActiveReplica);
140-
}
141-
142137
if (maxPartitionsPerInstance >= 0) {
143138
_record.setIntField(ResourceConfigProperty.MAX_PARTITIONS_PER_INSTANCE.name(), maxPartitionsPerInstance);
144139
}
@@ -227,15 +222,6 @@ public String getStateModelFactoryName() {
227222
return _record.getSimpleField(ResourceConfigProperty.STATE_MODEL_FACTORY_NAME.name());
228223
}
229224

230-
/**
231-
* Get the number of minimal active partitions for this resource.
232-
*
233-
* @return
234-
*/
235-
public int getMinActiveReplica() {
236-
return _record.getIntField(ResourceConfigProperty.MIN_ACTIVE_REPLICAS.name(), -1);
237-
}
238-
239225
// Delimiter used for storing active states as a comma-separated string in simpleField
240226
private static final String ACTIVE_STATES_DELIMITER = ",";
241227

@@ -574,7 +560,6 @@ public static class Builder {
574560
private Boolean _monitorDisabled;
575561
private int _numPartitions;
576562
private String _stateModelFactoryName;
577-
private int _minActiveReplica = -1;
578563
private int _maxPartitionsPerInstance = -1;
579564
private String _instanceGroupTag;
580565
private Boolean _helixEnabled;
@@ -632,15 +617,6 @@ public Builder setStateModelFactoryName(String stateModelFactoryName) {
632617
return this;
633618
}
634619

635-
public int getMinActiveReplica() {
636-
return _minActiveReplica;
637-
}
638-
639-
public Builder setMinActiveReplica(int minActiveReplica) {
640-
_minActiveReplica = minActiveReplica;
641-
return this;
642-
}
643-
644620
public int getMaxPartitionsPerInstance() {
645621
return _maxPartitionsPerInstance;
646622
}
@@ -789,7 +765,7 @@ public ResourceConfig build() {
789765
// validate();
790766

791767
return new ResourceConfig(_resourceId, _monitorDisabled, _numPartitions,
792-
_stateModelFactoryName, _minActiveReplica, _maxPartitionsPerInstance,
768+
_stateModelFactoryName, _maxPartitionsPerInstance,
793769
_instanceGroupTag, _helixEnabled, _externalViewDisabled, _rebalanceConfig,
794770
_stateTransitionTimeoutConfig, _preferenceLists, _mapFields, _p2pMessageEnabled,
795771
_partitionCapacityMap);
@@ -837,9 +813,6 @@ public static ResourceConfig mergeIdealStateWithResourceConfig(
837813
mergedZNRecord.setSimpleFieldIfAbsent(
838814
ResourceConfig.ResourceConfigProperty.STATE_MODEL_FACTORY_NAME.name(),
839815
idealState.getStateModelFactoryName());
840-
mergedZNRecord
841-
.setIntFieldIfAbsent(ResourceConfig.ResourceConfigProperty.MIN_ACTIVE_REPLICAS.name(),
842-
idealState.getMinActiveReplicas());
843816
mergedZNRecord
844817
.setBooleanFieldIfAbsent(ResourceConfig.ResourceConfigProperty.HELIX_ENABLED.name(),
845818
idealState.isEnabled());

helix-core/src/test/java/org/apache/helix/controller/rebalancer/waged/model/TestClusterModelProvider.java

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -308,7 +308,6 @@ private void prepareData(Map<String, Map<String, Map<String, String>>> input,
308308
Map<String, IdealState> isMap = new HashMap<>();
309309
for (String resource : _resourceNames) {
310310
ResourceConfig resourceConfig = new ResourceConfig.Builder(resource)
311-
.setMinActiveReplica(minActiveReplica)
312311
.build();
313312
_resourceConfigMap.put(resource, resourceConfig);
314313
IdealState is = new IdealState(resource);

helix-core/src/test/java/org/apache/helix/model/TestResourceConfig.java

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -187,10 +187,9 @@ public void testWithResourceBuilderInvalidInput() {
187187
@Test
188188
public void testConstructorWithoutGroupRoutingFields() {
189189
ResourceConfig resourceConfig = new ResourceConfig("resource", false, 2, "DEFAULT",
190-
1, 10, "placementTag", true, false, null, null, null, null, true);
190+
10, "placementTag", true, false, null, null, null, null, true);
191191
Assert.assertEquals(resourceConfig.getNumPartitions(), 2);
192192
Assert.assertEquals(resourceConfig.getStateModelFactoryName(), "DEFAULT");
193-
Assert.assertEquals(resourceConfig.getMinActiveReplica(), 1);
194193
Assert.assertEquals(resourceConfig.getMaxPartitionsPerInstance(), 10);
195194
Assert.assertEquals(resourceConfig.getInstanceGroupTag(), "placementTag");
196195
Assert.assertTrue(resourceConfig.isEnabled());
@@ -199,7 +198,7 @@ public void testConstructorWithoutGroupRoutingFields() {
199198
Assert.assertTrue(resourceConfig.isP2PMessageEnabled());
200199
for (String legacyField :
201200
new String[]{"RESOURCE_GROUP_NAME", "RESOURCE_TYPE", "GROUP_ROUTING_ENABLED",
202-
"STATE_MODEL_DEF_REF", "REPLICAS"}) {
201+
"STATE_MODEL_DEF_REF", "REPLICAS", "MIN_ACTIVE_REPLICAS"}) {
203202
Assert.assertFalse(resourceConfig.getRecord().getSimpleFields().containsKey(legacyField));
204203
}
205204
}
@@ -239,8 +238,6 @@ public void testMergeWithIdealState() {
239238
Assert.assertEquals(mergedResourceConfig.getNumPartitions(), testIdealState.getNumPartitions());
240239
Assert.assertEquals(mergedResourceConfig.getStateModelFactoryName(),
241240
testIdealState.getStateModelFactoryName());
242-
Assert.assertEquals(mergedResourceConfig.getMinActiveReplica(),
243-
testIdealState.getMinActiveReplicas());
244241
Assert
245242
.assertEquals(mergedResourceConfig.isEnabled().booleanValue(), testIdealState.isEnabled());
246243
for (String legacyField :
@@ -258,7 +255,6 @@ public void testMergeWithIdealState() {
258255
configBuilder.setMaxPartitionsPerInstance(2);
259256
configBuilder.setNumPartitions(2);
260257
configBuilder.setStateModelFactoryName("testRCFactory");
261-
configBuilder.setMinActiveReplica(2);
262258
configBuilder.setHelixEnabled(false);
263259
configBuilder.setExternalViewDisabled(true);
264260
testConfig = configBuilder.build();
@@ -274,8 +270,6 @@ public void testMergeWithIdealState() {
274270
Assert.assertEquals(mergedResourceConfig.getNumPartitions(), testConfig.getNumPartitions());
275271
Assert.assertEquals(mergedResourceConfig.getStateModelFactoryName(),
276272
testConfig.getStateModelFactoryName());
277-
Assert
278-
.assertEquals(mergedResourceConfig.getMinActiveReplica(), testConfig.getMinActiveReplica());
279273
Assert.assertEquals(mergedResourceConfig.isEnabled(), testConfig.isEnabled());
280274
for (String legacyField :
281275
new String[]{"RESOURCE_GROUP_NAME", "RESOURCE_TYPE", "GROUP_ROUTING_ENABLED"}) {

0 commit comments

Comments
 (0)