This is an automated email from the ASF dual-hosted git repository.
xiangfu0 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 ee615670b3f Reduce per-segment allocations in broker routing (#19530)
ee615670b3f is described below
commit ee615670b3f69a91eaba505005581f17244544da
Author: Xiang Fu <[email protected]>
AuthorDate: Tue Sep 15 18:49:41 2026 -0700
Reduce per-segment allocations in broker routing (#19530)
Reduce broker routing allocations using primitive pool counters and flat
required-segment maps with direct traversal. Canonicalize instance-config IDs
at refresh time and preserve routing behavior with regression tests.
---
.../instanceselector/BalancedInstanceSelector.java | 16 ++--
.../ReplicaGroupInstanceSelector.java | 15 +--
.../routing/manager/BaseBrokerRoutingManager.java | 19 ++--
.../instanceselector/InstanceSelectorTest.java | 62 +++++++++++++
.../routing/manager/BrokerRoutingManagerTest.java | 103 +++++++++++++++++++++
5 files changed, 194 insertions(+), 21 deletions(-)
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/BalancedInstanceSelector.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/BalancedInstanceSelector.java
index a1a7ce09bf2..cbb96b410b9 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/BalancedInstanceSelector.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/BalancedInstanceSelector.java
@@ -18,13 +18,15 @@
*/
package org.apache.pinot.broker.routing.instanceselector;
+import it.unimi.dsi.fastutil.ints.Int2IntMap;
+import it.unimi.dsi.fastutil.ints.Int2IntOpenHashMap;
+import it.unimi.dsi.fastutil.objects.Object2ObjectOpenHashMap;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import
org.apache.pinot.broker.routing.adaptiveserverselector.ServerSelectionContext;
import org.apache.pinot.common.metrics.BrokerMeter;
import org.apache.pinot.common.metrics.BrokerMetrics;
-import org.apache.pinot.common.utils.HashUtil;
/// Instance selector to balance the number of segments served by each
selected server instance.
///
@@ -45,11 +47,11 @@ public class BalancedInstanceSelector extends
BaseInstanceSelector {
@Override
public InstanceMapping select(List<String> segments, int requestId,
SegmentStates segmentStates, Map<String, String> queryOptions) {
- Map<String, String> segmentToSelectedInstanceMap = new
HashMap<>(HashUtil.getHashMapCapacity(segments.size()));
+ Map<String, String> segmentToSelectedInstanceMap = new
Object2ObjectOpenHashMap<>(segments.size());
// No need to adjust this map per total segment numbers, as optional
segments should be empty most of the time.
Map<String, String> optionalSegmentToInstanceMap = new HashMap<>();
ServerSelectionContext ctx = new ServerSelectionContext(queryOptions,
_config);
- Map<Integer, Integer> poolToSegmentCount = new HashMap<>();
+ Int2IntOpenHashMap poolToSegmentCount = new Int2IntOpenHashMap(2);
for (String segment : segments) {
List<SegmentInstanceCandidate> candidates =
segmentStates.getCandidates(segment);
@@ -72,7 +74,7 @@ public class BalancedInstanceSelector extends
BaseInstanceSelector {
} else {
selectedCandidate = candidates.get(requestId++ % candidates.size());
}
- poolToSegmentCount.merge(selectedCandidate.getPool(), 1, Integer::sum);
+ poolToSegmentCount.addTo(selectedCandidate.getPool(), 1);
// This can only be offline when it is a new segment. And such segment
is marked as optional segment so that
// broker or server can skip it upon any issue to process it.
if (selectedCandidate.isOnline()) {
@@ -82,9 +84,9 @@ public class BalancedInstanceSelector extends
BaseInstanceSelector {
}
}
- for (Map.Entry<Integer, Integer> entry : poolToSegmentCount.entrySet()) {
- _brokerMetrics.addMeteredValue(BrokerMeter.POOL_SEG_QUERIES,
entry.getValue(),
- BrokerMetrics.getTagForPreferredPool(queryOptions),
String.valueOf(entry.getKey()));
+ for (Int2IntMap.Entry entry : poolToSegmentCount.int2IntEntrySet()) {
+ _brokerMetrics.addMeteredValue(BrokerMeter.POOL_SEG_QUERIES,
entry.getIntValue(),
+ BrokerMetrics.getTagForPreferredPool(queryOptions),
String.valueOf(entry.getIntKey()));
}
return new InstanceMapping(segmentToSelectedInstanceMap,
optionalSegmentToInstanceMap);
}
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/ReplicaGroupInstanceSelector.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/ReplicaGroupInstanceSelector.java
index 493e9555f82..038994d40cd 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/ReplicaGroupInstanceSelector.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/routing/instanceselector/ReplicaGroupInstanceSelector.java
@@ -18,6 +18,9 @@
*/
package org.apache.pinot.broker.routing.instanceselector;
+import it.unimi.dsi.fastutil.ints.Int2IntMap;
+import it.unimi.dsi.fastutil.ints.Int2IntOpenHashMap;
+import it.unimi.dsi.fastutil.objects.Object2ObjectOpenHashMap;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
@@ -99,10 +102,10 @@ public class ReplicaGroupInstanceSelector extends
BaseInstanceSelector {
protected InstanceMapping selectServers(List<String> segments, int requestId,
SegmentStates segmentStates, @Nullable Map<String, Integer>
serverRankMap, ServerSelectionContext ctx) {
- Map<String, String> segmentToSelectedInstanceMap = new
HashMap<>(HashUtil.getHashMapCapacity(segments.size()));
+ Map<String, String> segmentToSelectedInstanceMap = new
Object2ObjectOpenHashMap<>(segments.size());
// No need to adjust this map per total segment numbers, as optional
segments should be empty most of the time.
Map<String, String> optionalSegmentToInstanceMap = new HashMap<>();
- Map<Integer, Integer> poolToSegmentCount = new HashMap<>();
+ Int2IntOpenHashMap poolToSegmentCount = new Int2IntOpenHashMap(2);
boolean useFixedReplica = ctx.isUseFixedReplica();
Integer numReplicaGroupsToQuery =
QueryOptionsUtils.getNumReplicaGroupsToQuery(ctx.getQueryOptions());
int numReplicaGroups = numReplicaGroupsToQuery != null ?
numReplicaGroupsToQuery : 1;
@@ -141,7 +144,7 @@ public class ReplicaGroupInstanceSelector extends
BaseInstanceSelector {
}
}
- poolToSegmentCount.merge(selectedInstance.getPool(), 1, Integer::sum);
+ poolToSegmentCount.addTo(selectedInstance.getPool(), 1);
// This can only be offline when it is a new segment. And such segment
is marked as optional segment so that
// broker or server can skip it upon any issue to process it.
if (selectedInstance.isOnline()) {
@@ -154,9 +157,9 @@ public class ReplicaGroupInstanceSelector extends
BaseInstanceSelector {
}
replicaOffset = (replicaOffset + 1) % numReplicaGroups;
}
- for (Map.Entry<Integer, Integer> entry : poolToSegmentCount.entrySet()) {
- _brokerMetrics.addMeteredValue(BrokerMeter.POOL_SEG_QUERIES,
entry.getValue(),
- BrokerMetrics.getTagForPreferredPool(ctx.getQueryOptions()),
String.valueOf(entry.getKey()));
+ for (Int2IntMap.Entry entry : poolToSegmentCount.int2IntEntrySet()) {
+ _brokerMetrics.addMeteredValue(BrokerMeter.POOL_SEG_QUERIES,
entry.getIntValue(),
+ BrokerMetrics.getTagForPreferredPool(ctx.getQueryOptions()),
String.valueOf(entry.getIntKey()));
}
return new InstanceMapping(segmentToSelectedInstanceMap,
optionalSegmentToInstanceMap);
}
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java
index db3f106b2e8..4b50b852bed 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/routing/manager/BaseBrokerRoutingManager.java
@@ -416,6 +416,8 @@ public abstract class BaseBrokerRoutingManager implements
RoutingManager, Cluste
String instanceId = instanceConfigZNRecord.getId();
try {
if (isEnabledServer(instanceConfigZNRecord)) {
+ // Join Jackson's JVM-interned IS/EV keys; a config-only Guava
interner would use a separate pool.
+ instanceId = instanceId.intern();
enabledServers.add(instanceId);
// Always refresh the server instance with the latest instance
config in case it changes
@@ -1193,30 +1195,31 @@ public abstract class BaseBrokerRoutingManager
implements RoutingManager, Cluste
private Map<ServerInstance, SegmentsToQuery>
getServerInstanceToSegmentsMap(String tableNameWithType,
InstanceSelector.SelectionResult selectionResult) {
Map<ServerInstance, SegmentsToQuery> merged = new HashMap<>();
- for (Map.Entry<String, String> entry :
selectionResult.getSegmentToInstanceMap().entrySet()) {
- ServerInstance serverInstance =
_enabledServerInstanceMap.get(entry.getValue());
+ // Flat selection maps can traverse their arrays directly without
allocating an entry object per segment.
+ selectionResult.getSegmentToInstanceMap().forEach((segment, instanceId) ->
{
+ ServerInstance serverInstance =
_enabledServerInstanceMap.get(instanceId);
if (serverInstance != null) {
SegmentsToQuery segmentsToQuery =
merged.computeIfAbsent(serverInstance, k -> new
SegmentsToQuery(new ArrayList<>(), new ArrayList<>()));
- segmentsToQuery.getSegments().add(entry.getKey());
+ segmentsToQuery.getSegments().add(segment);
} else {
// Should not happen in normal case unless encountered unexpected
exception when updating routing entries
_brokerMetrics.addMeteredTableValue(tableNameWithType,
BrokerMeter.SERVER_MISSING_FOR_ROUTING, 1L);
}
- }
- for (Map.Entry<String, String> entry :
selectionResult.getOptionalSegmentToInstanceMap().entrySet()) {
- ServerInstance serverInstance =
_enabledServerInstanceMap.get(entry.getValue());
+ });
+ selectionResult.getOptionalSegmentToInstanceMap().forEach((segment,
instanceId) -> {
+ ServerInstance serverInstance =
_enabledServerInstanceMap.get(instanceId);
if (serverInstance != null) {
SegmentsToQuery segmentsToQuery = merged.get(serverInstance);
// Skip servers that don't have non-optional segments, so that servers
always get some non-optional segments
// to process, to be backward compatible.
// TODO: allow servers only with optional segments
if (segmentsToQuery != null) {
- segmentsToQuery.getOptionalSegments().add(entry.getKey());
+ segmentsToQuery.getOptionalSegments().add(segment);
}
}
// TODO: Report missing server metrics when we allow servers only with
optional segments.
- }
+ });
return merged;
}
diff --git
a/pinot-broker/src/test/java/org/apache/pinot/broker/routing/instanceselector/InstanceSelectorTest.java
b/pinot-broker/src/test/java/org/apache/pinot/broker/routing/instanceselector/InstanceSelectorTest.java
index a7fe5b32e10..64bed8992ea 100644
---
a/pinot-broker/src/test/java/org/apache/pinot/broker/routing/instanceselector/InstanceSelectorTest.java
+++
b/pinot-broker/src/test/java/org/apache/pinot/broker/routing/instanceselector/InstanceSelectorTest.java
@@ -45,6 +45,7 @@ import
org.apache.pinot.broker.routing.adaptiveserverselector.HybridSelector;
import org.apache.pinot.common.metadata.ZKMetadataProvider;
import org.apache.pinot.common.metadata.segment.SegmentZKMetadata;
import org.apache.pinot.common.metrics.BrokerGauge;
+import org.apache.pinot.common.metrics.BrokerMeter;
import org.apache.pinot.common.metrics.BrokerMetrics;
import org.apache.pinot.common.request.BrokerRequest;
import org.apache.pinot.common.request.PinotQuery;
@@ -63,6 +64,7 @@ import org.testng.annotations.Test;
import static
org.apache.pinot.spi.config.table.RoutingConfig.REPLICA_GROUP_INSTANCE_SELECTOR_TYPE;
import static
org.apache.pinot.spi.config.table.RoutingConfig.STRICT_REPLICA_GROUP_INSTANCE_SELECTOR_TYPE;
+import static
org.apache.pinot.spi.utils.CommonConstants.Broker.Request.QueryOptionKey.ORDERED_PREFERRED_POOLS;
import static
org.apache.pinot.spi.utils.CommonConstants.Helix.StateModel.SegmentStateModel.CONSUMING;
import static
org.apache.pinot.spi.utils.CommonConstants.Helix.StateModel.SegmentStateModel.ERROR;
import static
org.apache.pinot.spi.utils.CommonConstants.Helix.StateModel.SegmentStateModel.OFFLINE;
@@ -76,6 +78,7 @@ import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
import static org.mockito.Mockito.when;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertNull;
@@ -2344,4 +2347,63 @@ public class InstanceSelectorTest {
verify(_brokerMetrics, never()).setValueOfTableGauge(eq(TABLE_NAME),
any(BrokerGauge.class), anyLong());
verify(_brokerMetrics, never()).removeTableGauge(eq(TABLE_NAME),
any(BrokerGauge.class));
}
+
+ @DataProvider(name = "poolMetricsSelector")
+ public Object[] getPoolMetricsSelector() {
+ return new Object[]{BALANCED_INSTANCE_SELECTOR,
REPLICA_GROUP_INSTANCE_SELECTOR_TYPE};
+ }
+
+ @Test(dataProvider = "poolMetricsSelector")
+ public void testSelectedPoolMetricsAndRequestIsolation(String selectorType) {
+ BaseInstanceSelector selector =
selectorType.equals(BALANCED_INSTANCE_SELECTOR)
+ ? new BalancedInstanceSelector()
+ : new ReplicaGroupInstanceSelector();
+ selector._brokerMetrics = _brokerMetrics;
+ selector._config = INSTANCE_SELECTOR_CONFIG;
+ Map<String, List<SegmentInstanceCandidate>> candidates = new HashMap<>();
+ Map<String, String> expectedInstances = new HashMap<>();
+ List<String> segments = new ArrayList<>();
+ // Exercise the fallback pool, the primitive-map zero key, and IDs/counts
outside the Integer cache.
+ for (int pool : new int[]{-1, 0, 128}) {
+ String instance = "instance" + pool;
+ for (int i = 0; i < 129; i++) {
+ String segment = "segment_" + pool + "_" + i;
+ candidates.put(segment, List.of(new SegmentInstanceCandidate(instance,
true, pool, 0)));
+ expectedInstances.put(segment, instance);
+ // Routing matches segment names by value, even when query and
metadata use different String objects.
+ segments.add(new String(segment));
+ }
+ }
+ candidates.put("optional", List.of(new
SegmentInstanceCandidate("instance128", false, 128, 0)));
+ segments.addAll(List.of("optional", "unavailable", "pending"));
+ selector._segmentStates = new SegmentStates(candidates,
Set.of("instance-1", "instance0", "instance128"),
+ Set.of("unavailable", "unrequestedUnavailable"));
+ Map<String, String> queryOptions = Map.of(ORDERED_PREFERRED_POOLS,
"128|0");
+ when(_pinotQuery.getQueryOptions()).thenReturn(queryOptions);
+
+ InstanceSelector.SelectionResult first = selector.select(_brokerRequest,
segments, 0L);
+
+ assertEquals(first.getSegmentToInstanceMap(), expectedInstances);
+ assertEquals(first.getOptionalSegmentToInstanceMap(), Map.of("optional",
"instance128"));
+ assertEquals(first.getUnavailableSegments(), List.of("unavailable"));
+ String preferredPoolTag =
BrokerMetrics.getTagForPreferredPool(queryOptions);
+ verify(_brokerMetrics).addMeteredValue(BrokerMeter.POOL_SEG_QUERIES, 129L,
preferredPoolTag, "-1");
+ verify(_brokerMetrics).addMeteredValue(BrokerMeter.POOL_SEG_QUERIES, 129L,
preferredPoolTag, "0");
+ // Optional segments count as selected; unavailable segments and metadata
not yet present do not.
+ verify(_brokerMetrics).addMeteredValue(BrokerMeter.POOL_SEG_QUERIES, 130L,
preferredPoolTag, "128");
+ verifyNoMoreInteractions(_brokerMetrics);
+
+ clearInvocations(_brokerMetrics);
+ when(_pinotQuery.getQueryOptions()).thenReturn(Map.of());
+ InstanceSelector.SelectionResult second = selector.select(_brokerRequest,
List.of("segment_0_0"), 1L);
+
+ assertEquals(second.getSegmentToInstanceMap(), Map.of("segment_0_0",
"instance0"));
+ assertTrue(second.getOptionalSegmentToInstanceMap().isEmpty());
+ assertTrue(second.getUnavailableSegments().isEmpty());
+ verify(_brokerMetrics).addMeteredValue(BrokerMeter.POOL_SEG_QUERIES, 1L,
+ BrokerMetrics.getTagForPreferredPool(Map.of()), "0");
+ verifyNoMoreInteractions(_brokerMetrics);
+ assertEquals(first.getSegmentToInstanceMap(), expectedInstances);
+ assertEquals(first.getOptionalSegmentToInstanceMap(), Map.of("optional",
"instance128"));
+ }
}
diff --git
a/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerTest.java
b/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerTest.java
index 75c2f3570ba..5d7830ad04b 100644
---
a/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerTest.java
+++
b/pinot-broker/src/test/java/org/apache/pinot/broker/routing/manager/BrokerRoutingManagerTest.java
@@ -18,6 +18,7 @@
*/
package org.apache.pinot.broker.routing.manager;
+import it.unimi.dsi.fastutil.objects.Object2ObjectOpenHashMap;
import java.lang.reflect.Constructor;
import java.util.HashSet;
import java.util.List;
@@ -35,6 +36,7 @@ import org.apache.helix.model.IdealState;
import org.apache.helix.model.InstanceConfig;
import org.apache.helix.store.zk.ZkHelixPropertyStore;
import org.apache.helix.zookeeper.datamodel.ZNRecord;
+import org.apache.helix.zookeeper.datamodel.serializer.ZNRecordSerializer;
import org.apache.pinot.broker.routing.instanceselector.InstanceSelector;
import org.apache.pinot.broker.routing.instanceselector.TableReplicaHealth;
import
org.apache.pinot.broker.routing.segmentmetadata.SegmentZkMetadataFetcher;
@@ -45,10 +47,13 @@ import
org.apache.pinot.broker.routing.segmentselector.SegmentSelector;
import org.apache.pinot.broker.routing.tablesampler.TableSampler;
import org.apache.pinot.broker.routing.timeboundary.TimeBoundaryManager;
import org.apache.pinot.common.metrics.BrokerGauge;
+import org.apache.pinot.common.metrics.BrokerMeter;
import org.apache.pinot.common.metrics.BrokerMetrics;
import org.apache.pinot.common.request.BrokerRequest;
import org.apache.pinot.common.request.QuerySource;
import org.apache.pinot.common.utils.config.TableConfigSerDeUtils;
+import org.apache.pinot.core.routing.RoutingTable;
+import org.apache.pinot.core.routing.SegmentsToQuery;
import org.apache.pinot.core.routing.TablePartitionInfo;
import org.apache.pinot.core.routing.TablePartitionReplicatedServersInfo;
import org.apache.pinot.core.routing.timeboundary.TimeBoundaryInfo;
@@ -57,6 +62,7 @@ import
org.apache.pinot.core.transport.server.routing.stats.ServerRoutingStatsMa
import org.apache.pinot.spi.config.table.TableConfig;
import org.apache.pinot.spi.config.table.TableType;
import org.apache.pinot.spi.env.PinotConfiguration;
+import org.apache.pinot.spi.utils.CommonConstants.Helix;
import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
import org.apache.pinot.spi.utils.builder.TableNameBuilder;
import org.apache.zookeeper.data.Stat;
@@ -71,15 +77,18 @@ import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.clearInvocations;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
import static org.mockito.Mockito.when;
import static org.testng.Assert.assertEquals;
import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNotSame;
import static org.testng.Assert.assertNull;
import static org.testng.Assert.assertSame;
import static org.testng.Assert.assertTrue;
@@ -159,6 +168,100 @@ public class BrokerRoutingManagerTest {
_mocks.close();
}
+ @Test
+ public void testInstanceConfigIdsInternedAcrossRefreshesAndRouting()
+ throws Exception {
+ ZNRecordSerializer serializer = new ZNRecordSerializer();
+ ZNRecord assignments = new ZNRecord(TEST_TABLE);
+ assignments.setMapField("required", Map.of(SERVER_INSTANCE_ID, "ONLINE"));
+ assignments.setMapField("optional", Map.of(SERVER_INSTANCE_ID, "ONLINE"));
+ byte[] serializedAssignments = serializer.serialize(assignments);
+ InstanceSelector instanceSelector = mock(InstanceSelector.class);
+ putRoutingEntry(TEST_TABLE,
+ createRoutingEntry(TEST_TABLE, selectorOf("required", "optional"),
List.of(), instanceSelector));
+
+ for (int grpcPort : List.of(9000, 9001)) {
+ ZNRecord decodedAssignments = (ZNRecord)
serializer.deserialize(serializedAssignments);
+ String requiredInstanceId =
decodedAssignments.getMapField("required").keySet().iterator().next();
+ String optionalInstanceId =
decodedAssignments.getMapField("optional").keySet().iterator().next();
+ // Exercise IDs from real assignment-map decoding, including the
parser's JVM interning behavior.
+ assertSame(requiredInstanceId, SERVER_INSTANCE_ID);
+ assertSame(optionalInstanceId, requiredInstanceId);
+ when(instanceSelector.select(any(), any(), anyLong())).thenReturn(new
InstanceSelector.SelectionResult(
+ new InstanceSelector.InstanceMapping(Map.of("required",
requiredInstanceId),
+ Map.of("optional", optionalInstanceId)), List.of(), 0));
+
+ ZNRecord config = createEnabledServerZNRecord(SERVER_INSTANCE_ID);
+ config.setIntField(Helix.Instance.GRPC_PORT_KEY, grpcPort);
+ ZNRecord decodedConfig = (ZNRecord)
serializer.deserialize(serializer.serialize(config));
+ // Config IDs are JSON values; unlike assignment-map keys, they are not
interned by the parser.
+ assertEquals(decodedConfig.getId(), SERVER_INSTANCE_ID);
+ assertNotSame(decodedConfig.getId(), requiredInstanceId);
+ when(_zkDataAccessor.getChildren(eq(INSTANCE_CONFIGS_PATH), any(),
eq(AccessOption.PERSISTENT), anyInt(),
+ anyInt())).thenReturn(List.of(decodedConfig));
+
+ _routingManager.processClusterChange(ChangeType.INSTANCE_CONFIG);
+
+ Map<String, ServerInstance> enabledServers =
_routingManager.getEnabledServerInstanceMap();
+ assertEquals(enabledServers.size(), 1);
+ assertSame(enabledServers.keySet().iterator().next(),
requiredInstanceId);
+ ServerInstance server = enabledServers.get(requiredInstanceId);
+ assertEquals(server.getInstanceId(), SERVER_INSTANCE_ID);
+ assertEquals(server.getHostname(), SERVER_HOST);
+ assertEquals(server.getPort(), SERVER_PORT);
+ // An equal ID on a later config refresh must still replace the server's
configuration.
+ assertEquals(server.getGrpcPort(), grpcPort);
+
assertSame(_routingManager.getRoutableServerInstanceMap().get(requiredInstanceId),
server);
+
+ RoutingTable routingTable =
_routingManager.getRoutingTable(brokerRequest(TEST_TABLE), 0);
+ Map<ServerInstance, SegmentsToQuery> serverSegments =
routingTable.getServerInstanceToSegmentsMap();
+ assertEquals(serverSegments.size(), 1);
+ assertSame(serverSegments.keySet().iterator().next(), server);
+ assertEquals(serverSegments.get(server).getSegments(),
List.of("required"));
+ assertEquals(serverSegments.get(server).getOptionalSegments(),
List.of("optional"));
+ assertTrue(routingTable.getUnavailableSegments().isEmpty());
+ }
+ }
+
+ @Test
+ public void testGroupingPreservesOptionalServerAndMissingServerBehavior()
+ throws Exception {
+ String optionalOnlyInstanceId = "Server_optional_8000";
+ ServerInstance requiredServer =
+ new ServerInstance(new
InstanceConfig(createEnabledServerZNRecord(SERVER_INSTANCE_ID)));
+ ServerInstance optionalOnlyServer =
+ new ServerInstance(new
InstanceConfig(createEnabledServerZNRecord(optionalOnlyInstanceId)));
+ _routingManager.getEnabledServerInstanceMap().put(SERVER_INSTANCE_ID,
requiredServer);
+ _routingManager.getEnabledServerInstanceMap().put(optionalOnlyInstanceId,
optionalOnlyServer);
+ Map<String, String> required = new Object2ObjectOpenHashMap<>(Map.of(
+ "required0", SERVER_INSTANCE_ID, "required1", SERVER_INSTANCE_ID,
+ "missingRequired0", "Server_missing_8000", "missingRequired1",
"Server_missing_8000"));
+ Map<String, String> optional = Map.of("optional", SERVER_INSTANCE_ID,
+ "optionalOnly", optionalOnlyInstanceId, "missingOptional",
"Server_missing_optional_8000");
+ InstanceSelector instanceSelector = mock(InstanceSelector.class);
+ when(instanceSelector.select(any(), any(), anyLong())).thenReturn(new
InstanceSelector.SelectionResult(
+ new InstanceSelector.InstanceMapping(required, optional), List.of(),
0));
+ putRoutingEntry(TEST_TABLE, createRoutingEntry(TEST_TABLE,
+ selectorOf("required0", "required1", "missingRequired0",
"missingRequired1", "optional", "optionalOnly",
+ "missingOptional"), List.of(), instanceSelector));
+ clearInvocations(_brokerMetrics);
+
+ RoutingTable routingTable =
_routingManager.getRoutingTable(brokerRequest(TEST_TABLE), 0);
+
+ Map<ServerInstance, SegmentsToQuery> grouped =
routingTable.getServerInstanceToSegmentsMap();
+ assertEquals(grouped.keySet(), Set.of(requiredServer));
+ assertEquals(new HashSet<>(grouped.get(requiredServer).getSegments()),
Set.of("required0", "required1"));
+ assertEquals(grouped.get(requiredServer).getSegments().size(), 2);
+ assertEquals(grouped.get(requiredServer).getOptionalSegments(),
List.of("optional"));
+ assertTrue(routingTable.getUnavailableSegments().isEmpty());
+ // Missing required segments are metered; missing optional segments and
optional-only servers are skipped.
+ ArgumentCaptor<Long> increments = ArgumentCaptor.forClass(Long.class);
+ verify(_brokerMetrics, atLeastOnce()).addMeteredTableValue(eq(TEST_TABLE),
+ eq(BrokerMeter.SERVER_MISSING_FOR_ROUTING), increments.capture());
+
assertEquals(increments.getAllValues().stream().mapToLong(Long::longValue).sum(),
2L);
+ verifyNoMoreInteractions(_brokerMetrics);
+ }
+
@Test
public void testNoErrorWhenCallbackNotSet() {
// Don't set callback
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]