This is an automated email from the ASF dual-hosted git repository.
Jackie-Jiang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/pinot.git
The following commit(s) were added to refs/heads/master by this push:
new 811d64c5923 Fix FD_AWARE instance assignment when pool tags do not
start at 0. (#19266)
811d64c5923 is described below
commit 811d64c5923c3578c976216a381e357b03d310fc
Author: deepinsight coder <[email protected]>
AuthorDate: Tue Aug 18 12:55:56 2026 -0700
Fix FD_AWARE instance assignment when pool tags do not start at 0. (#19266)
---
.../instance/FDAwareInstancePartitionSelector.java | 6 --
.../instance/InstanceAssignmentTest.java | 101 +++++++++++++++++++++
2 files changed, 101 insertions(+), 6 deletions(-)
diff --git
a/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/assignment/instance/FDAwareInstancePartitionSelector.java
b/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/assignment/instance/FDAwareInstancePartitionSelector.java
index 6a297df1f4b..f909e3ac643 100644
---
a/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/assignment/instance/FDAwareInstancePartitionSelector.java
+++
b/pinot-controller/src/main/java/org/apache/pinot/controller/helix/core/assignment/instance/FDAwareInstancePartitionSelector.java
@@ -286,7 +286,6 @@ public class FDAwareInstancePartitionSelector extends
InstancePartitionSelector
int _mapDimInstancePerReplicaGroup;
HashMap<String, Integer> _usedInstances = new HashMap<>();
int _numFaultDomains;
- int[][] _fdCounter;
ReplicaGroupBasedAssignmentState(int numReplicaGroups, int
numInstancesPerReplicaGroup,
int numExistingReplicaGroups, int numExistingInstancesPerReplicaGroup,
int numFaultDomains) {
@@ -300,7 +299,6 @@ public class FDAwareInstancePartitionSelector extends
InstancePartitionSelector
_mapDimInstancePerReplicaGroup =
Math.max(numExistingInstancesPerReplicaGroup, numInstancesPerReplicaGroup);
_replicaGroupIdToInstancesMap = new
Instance[_mapDimReplicaGroup][_mapDimInstancePerReplicaGroup];
- _fdCounter = new int[_mapDimInstancePerReplicaGroup][_numFaultDomains];
}
ReplicaGroupBasedAssignmentState(int numReplicaGroups, int
numInstancesPerReplicaGroup, int numFaultDomains) {
@@ -324,20 +322,16 @@ public class FDAwareInstancePartitionSelector extends
InstancePartitionSelector
Preconditions.checkState(instance.getExistingReplicaGroupId() ==
Instance.NEW_INSTANCE);
_replicaGroupIdToInstancesMap[replicaGroupId][instanceIndex] = instance;
_usedInstances.put(instance.getInstanceName(),
instance.getFaultDomainId());
- _fdCounter[instanceIndex][instance.getFaultDomainId()] += 1;
}
private void setExistingInstance(int replicaGroupId, int instanceIndex,
String instance, int fdId) {
_replicaGroupIdToInstancesMap[replicaGroupId][instanceIndex] = new
Instance(instance, fdId, replicaGroupId);
_usedInstances.put(instance, fdId);
- _fdCounter[instanceIndex][fdId] += 1;
}
private void unSetInstance(int replicaGroupId, int instanceIndex) {
- int fdId =
_replicaGroupIdToInstancesMap[replicaGroupId][instanceIndex].getFaultDomainId();
_usedInstances.remove(_replicaGroupIdToInstancesMap[replicaGroupId][instanceIndex].getInstanceName());
_replicaGroupIdToInstancesMap[replicaGroupId][instanceIndex] = null;
- _fdCounter[instanceIndex][fdId] -= 1;
}
/// From an exising replica group, remove the instances that are gone and
set them in \_replicaGroupIdToInstancesMap
diff --git
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/assignment/instance/InstanceAssignmentTest.java
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/assignment/instance/InstanceAssignmentTest.java
index e8abb617079..153fb05055b 100644
---
a/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/assignment/instance/InstanceAssignmentTest.java
+++
b/pinot-controller/src/test/java/org/apache/pinot/controller/helix/core/assignment/instance/InstanceAssignmentTest.java
@@ -3393,6 +3393,107 @@ public class InstanceAssignmentTest {
assertEquals(steadyStatePartitions.getInstances(0, rg),
initialPartitions.getInstances(0, rg));
}
}
+
+ /// FD_AWARE must accept Helix pool tags that do not start at 0 (e.g. 1/2/3).
+ /// Every other pool-based test uses {@code pool = i % numPools}, so this
crash had no coverage.
+ /// See <a href="https://github.com/apache/pinot/issues/12239">#12239</a>.
+ @Test
+ public void testPoolBasedFDAwareNonZeroBasedPools() {
+ // Same first topology as testPoolBasedFDAware, but pool = (i % numPools)
+ 1 → {1,2,3,4,5}.
+ int numInstances = 21;
+ int numPools = 5;
+ int numReplicaGroups = 3;
+ int numInstancesPerReplicaGroup = numInstances / numReplicaGroups;
+ List<InstanceConfig> instanceConfigs =
+ newFDAwarePoolInstanceConfigs(numInstances, numPools);
+ InstanceTagPoolConfig tagPoolConfig = new
InstanceTagPoolConfig(OFFLINE_TAG, true, numPools, null);
+ InstanceReplicaGroupPartitionConfig replicaPartitionConfig =
+ new InstanceReplicaGroupPartitionConfig(true, 0, numReplicaGroups,
numInstancesPerReplicaGroup, 0, 0, false,
+ null);
+ TableConfig tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME)
+
.setInstanceAssignmentConfigMap(Map.of(InstancePartitionsType.OFFLINE.toString(),
+ new InstanceAssignmentConfig(tagPoolConfig, null,
replicaPartitionConfig,
+
InstanceAssignmentConfig.PartitionSelector.FD_AWARE_INSTANCE_PARTITION_SELECTOR.toString(),
false)))
+ .build();
+ InstanceAssignmentDriver driver = new
InstanceAssignmentDriver(tableConfig);
+ InstancePartitions instancePartitions =
+ driver.assignInstances(InstancePartitionsType.OFFLINE,
instanceConfigs, null);
+ assertFilledUniqueAssignment(instancePartitions, instanceConfigs,
numReplicaGroups,
+ numInstancesPerReplicaGroup);
+
+ // Incremental uplift with minimizeDataMovement hits setExistingInstance
with the raw pool ids.
+ numInstances = 28;
+ numReplicaGroups = 4;
+ numInstancesPerReplicaGroup = numInstances / numReplicaGroups;
+ instanceConfigs = newFDAwarePoolInstanceConfigs(numInstances, numPools);
+ replicaPartitionConfig =
+ new InstanceReplicaGroupPartitionConfig(true, 0, numReplicaGroups,
numInstancesPerReplicaGroup, 0, 0, true,
+ null);
+ tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME)
+
.setInstanceAssignmentConfigMap(Map.of(InstancePartitionsType.OFFLINE.toString(),
+ new InstanceAssignmentConfig(tagPoolConfig, null,
replicaPartitionConfig,
+
InstanceAssignmentConfig.PartitionSelector.FD_AWARE_INSTANCE_PARTITION_SELECTOR.toString(),
true)))
+ .build();
+ driver = new InstanceAssignmentDriver(tableConfig);
+ instancePartitions =
driver.assignInstances(InstancePartitionsType.OFFLINE, instanceConfigs,
instancePartitions);
+ assertFilledUniqueAssignment(instancePartitions, instanceConfigs,
numReplicaGroups,
+ numInstancesPerReplicaGroup);
+
+ // Reporter case: pools {1,2,3}.
+ numInstances = 6;
+ numPools = 3;
+ numReplicaGroups = 3;
+ numInstancesPerReplicaGroup = 2;
+ instanceConfigs = newFDAwarePoolInstanceConfigs(numInstances, numPools);
+ tagPoolConfig = new InstanceTagPoolConfig(OFFLINE_TAG, true, numPools,
null);
+ replicaPartitionConfig =
+ new InstanceReplicaGroupPartitionConfig(true, 0, numReplicaGroups,
numInstancesPerReplicaGroup, 0, 0, false,
+ null);
+ tableConfig = new
TableConfigBuilder(TableType.OFFLINE).setTableName(RAW_TABLE_NAME)
+
.setInstanceAssignmentConfigMap(Map.of(InstancePartitionsType.OFFLINE.toString(),
+ new InstanceAssignmentConfig(tagPoolConfig, null,
replicaPartitionConfig,
+
InstanceAssignmentConfig.PartitionSelector.FD_AWARE_INSTANCE_PARTITION_SELECTOR.toString(),
false)))
+ .build();
+ driver = new InstanceAssignmentDriver(tableConfig);
+ instancePartitions =
driver.assignInstances(InstancePartitionsType.OFFLINE, instanceConfigs, null);
+ assertFilledUniqueAssignment(instancePartitions, instanceConfigs,
numReplicaGroups,
+ numInstancesPerReplicaGroup);
+ }
+
+ /// Builds pool-tagged instance configs whose pool ids start at 1 rather
than 0.
+ private static List<InstanceConfig> newFDAwarePoolInstanceConfigs(int
numInstances, int numPools) {
+ List<InstanceConfig> instanceConfigs = new ArrayList<>(numInstances);
+ for (int i = 0; i < numInstances; i++) {
+ int pool = (i % numPools) + 1;
+ InstanceConfig instanceConfig =
+ new InstanceConfig(SERVER_INSTANCE_ID_PREFIX + i +
SERVER_INSTANCE_POOL_PREFIX + pool);
+ instanceConfig.addTag(OFFLINE_TAG);
+ instanceConfig.getRecord().setMapField(InstanceUtils.POOL_KEY,
Map.of(OFFLINE_TAG, Integer.toString(pool)));
+ instanceConfigs.add(instanceConfig);
+ }
+ return instanceConfigs;
+ }
+
+ private static void assertFilledUniqueAssignment(InstancePartitions
instancePartitions,
+ List<InstanceConfig> instanceConfigs, int numReplicaGroups, int
numInstancesPerReplicaGroup) {
+ assertEquals(instancePartitions.getNumReplicaGroups(), numReplicaGroups);
+ assertEquals(instancePartitions.getNumPartitions(), 1);
+ Set<String> candidates = new HashSet<>();
+ for (InstanceConfig instanceConfig : instanceConfigs) {
+ candidates.add(instanceConfig.getInstanceName());
+ }
+ Set<String> assigned = new HashSet<>();
+ for (int replicaGroupId = 0; replicaGroupId < numReplicaGroups;
replicaGroupId++) {
+ List<String> instances = instancePartitions.getInstances(0,
replicaGroupId);
+ assertEquals(instances.size(), numInstancesPerReplicaGroup);
+ for (String instance : instances) {
+ assertTrue(candidates.contains(instance), instance);
+ assertTrue(assigned.add(instance), instance);
+ }
+ }
+ assertEquals(assigned.size(), numReplicaGroups *
numInstancesPerReplicaGroup);
+ }
+
/// Verifies that subset-partition tables use the total Kafka partition
count (not the subset size)
/// for instance assignment, producing the same server spread as a normal
full-partition table.
///
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]