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 3ae199f90c6 [perf] Encode MSE segment lists at dispatch, optionally as
protobuf (#19568)
3ae199f90c6 is described below
commit 3ae199f90c663e661ae80c02edefade2cf66366c
Author: Xiang Fu <[email protected]>
AuthorDate: Wed Sep 23 17:02:20 2026 -0700
[perf] Encode MSE segment lists at dispatch, optionally as protobuf (#19568)
* Stop JSON-encoding MSE leaf-stage segment lists on the broker compile path
On tables with many segments, the top broker CPU consumer of a multi-stage
query was `WorkerMetadata.setTableSegmentsMap`: for every leaf-stage worker
the planner JSON-encoded the full list of routed segment names with Jackson,
on the fixed-size `multi-stage-query-compile-executor`, only for the
dispatcher to copy that string into a proto `map<string, string>` custom
property and for every server to JSON-parse it back (once per worker).
- `WorkerMetadata` now holds the segment maps as plain objects; nothing is
encoded at plan time anymore.
- `QueryPlanSerDeUtils` encodes them per server at dispatch time, once per
worker, in one of two wire encodings picked per request:
- proto: new native `tableSegmentsMap` / `logicalTableSegmentsMap` fields
of `Worker.WorkerMetadata` (`SegmentsMap` / `SegmentList` messages).
- legacy JSON: the previous custom-property string, still the default.
Decoding accepts both, so a new server understands every broker.
- New broker config `pinot.broker.mse.proto.segment.list` (default false)
and query option `protoSegmentList` to enable the proto encoding. It has
to stay opt-in for one release: an older server finds no segments under
the proto fields and fails the leaf stage, so it may only be enabled once
the whole fleet is upgraded.
- That config is live: `ProtoSegmentListPredicate` seeds it from the static
broker config and then follows cluster config on the same key, registered
by `BaseBrokerStarter` against the existing cluster-config change handler,
and `QueryDispatcher` reads it per request. The precondition is exactly
"every server already upgraded", so operators need to turn it on the
moment a rolling upgrade completes, and to turn it back off at once if it
misbehaves; neither should cost a broker restart. Precedence is per-query
`SET protoSegmentList`, then cluster config, then static broker config,
and clearing the cluster-config key falls back to the shipped `false`
rather than to the static seed.
Micro-benchmark, one leaf worker, 60-char segment names (Java 25, x86):
segments | legacy JSON encode | proto encode | legacy decode | proto
decode
1,000 | 177 us | 132 us | 66 us | 29
us
3,000 | 624 us | 394 us | 200 us | 95
us
20,000 | 5,101 us | 2,709 us | 1,419 us | 782
us
60,000 | 15,950 us | 8,217 us | 4,602 us | 2,386
us
* Auto-enable the proto segment list encoding from server versions
Address review feedback on the proto segment list encoding.
- Replace the live cluster-config toggle with NEVER / SAFE / ALWAYS modes
on `pinot.broker.mse.proto.segment.list`, modeled on SendStatsPredicate.
SAFE (the default) watches the Helix instance configs and uses the proto
encoding only while every server reports this broker's Pinot version, so
it switches itself on when a rolling upgrade completes. It fails closed:
missing, UNKNOWN or unreadable versions count as outdated, it stays off
until the first delivery, and multi-cluster queries always use the
legacy encoding because remote-cluster servers are not watched.
- Keep QueryDispatcher's 10-arg constructor and the 1-arg
QueryPlanSerDeUtils.toProtoWorkerMetadataList as overloads for
out-of-tree callers; both keep the legacy encoding.
- Keep decoded worker custom properties unmodifiable under both encodings.
- Drop segment maps from EXPLAIN responses; the broker reads only the root
node of each stage plan.
- Document the no-mutation contract of the segment map setters, the
plan-time invariant in DispatchablePlanContext, and reword the stage-0
comment in PinotDispatchPlanner.
- Tests: SAFE mode logic and Helix translation, zero-entry segments maps
through real bytes, unmodifiable custom properties, and integration tests
comparing both encodings for physical and logical tables, including SAFE
switching itself on in a real cluster.
* Make the segment list encoding mode a live cluster config
The per-query `protoSegmentList` option was the only way to change the
encoding without restarting a broker, but asking clients to change their
queries is harder than restarting brokers, so it was the wrong escape
hatch. Read the mode from cluster config instead, and drop the option.
- `ProtoSegmentListPredicate` also listens on cluster config for
`pinot.broker.mse.proto.segment.list`. Precedence is cluster config,
then the static broker config, then SAFE; clearing the cluster-config
key restores the static broker config, as MultiStageQueryThrottler does
for its live config. A value that is not a mode is ignored with a
warning rather than moving the cluster off the operator's chosen mode.
- The server versions are watched whatever the static mode is, since
cluster config can select SAFE at runtime.
- Remove the `protoSegmentList` query option, `QueryOptionsUtils
.isProtoSegmentList` and its test. Nothing overrides the mode per query
any more, so all servers of a query always agree on the encoding.
- The integration tests now switch the encoding through cluster config
rather than a query option, which also covers the live reload path: the
logical-table test goes through the controller's /cluster/configs
endpoint and restores SAFE afterwards, since that cluster is shared.
* Make the proto segment list encoding a plain live cluster config
The version-based auto-detection was more machinery than the problem
needs. The encoding ships disabled, so a rolling upgrade is safe in any
broker/server order without inferring anything, and switching it on once
the fleet is uniform is a deliberate operator action.
- `pinot.broker.mse.proto.segment.list` is a boolean again, default
false, read from cluster config as well as the static broker config.
Cluster config wins and takes effect on the next query, so it can be
turned on, and off again, without restarting the brokers. Clearing the
key restores the static broker config, and a value that is neither
true nor false reads as disabled.
- `QueryDispatcher` holds the flag and is the cluster-config listener, so
`ProtoSegmentListPredicate`, its modes, its instance-config watching and
the `getProtoSegmentListPredicate()` seam on the request handler all go
away. Enabling logs a WARN naming the precondition: every server the
broker dispatches to, including remote clusters, must understand the
proto fields.
- The multi-cluster special case goes with it: the precondition covers
remote-cluster servers, rather than the broker silently opting those
queries out.
- Tests: `QueryDispatcherTest` covers the live flag, its precedence over
the static config and its fail-closed parsing; the integration tests
turn the encoding on through cluster config and compare both encodings.
* Rename the proto segment list config to say that it is a toggle
DEFAULT_MSE_PROTO_SEGMENT_LIST read as a noun, as though the constant
held a segment list rather than whether the encoding is on. Rename to
CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST /
DEFAULT_MSE_ENABLE_PROTO_SEGMENT_LIST
with the key pinot.broker.mse.enable.proto.segment.list, matching the
sibling boolean CONFIG_OF_MSE_ENABLE_GROUP_TRIM
(pinot.broker.mse.enable.group.trim).
* Fall back to disabled when the cluster-config key is cleared
Tracking the static broker config as a separate fallback value was extra
state for no benefit: clearing the key now simply disables the encoding,
which is the old behaviour every server understands. Drops the
_staticProtoSegmentList field, collapses the parsing to one expression
where anything other than true reads as disabled, and shortens the
repeated config key reference to a local constant.
* Name the dispatcher flag after the config it follows
Rename _protoSegmentList to _enableProtoSegmentList, with the matching
local constant, parameters, locals and isEnableProtoSegmentList()
accessor, so the field reads the same way as
pinot.broker.mse.enable.proto.segment.list. The serde parameter in
QueryPlanSerDeUtils keeps its name: it selects an encoding rather than
following the config.
* Add a JMH benchmark for the segment list encodings
BenchmarkSegmentListEncoding measures both encodings of one leaf-stage
worker through the real QueryPlanSerDeUtils paths: building the request
on the broker, building plus serializing it, and decoding it on a server.
It lives in the org.apache.pinot.query.routing package so that it can
call the package-private decode entry point the server uses, rather than
widening that method for a benchmark.
The numbers it produces replace the ones quoted in the PR description,
which overstated the dispatch-side win: building the proto form is much
cheaper than building the JSON string, but serializing 60k individually
length-prefixed strings costs about as much as one large JSON string.
* Parse legacy segment maps lazily, on each worker's own thread
Addresses review feedback: decoding the legacy JSON segment maps eagerly
in fromProtoWorkerMetadata made the server parse every worker's list one
after another in deserializePlan, before any worker of the stage started.
On master each worker parsed its own list lazily, when its leaf stage was
compiled, so the parses ran in parallel on the workers' threads. Since
the proto encoding ships disabled, this was the default path.
- WorkerMetadata keeps a legacy-encoded segment map as the raw JSON and
parses it on first access, memoized through a volatile field (a racing
first access at worst parses twice). isLeafStageWorker() counts an
unparsed map. The JSON setters are package-private, for
QueryPlanSerDeUtils only.
- The proto path is unchanged: it stays eager, since protobuf has already
materialized the strings while parsing the request.
- Malformed JSON now fails at the worker's first access again, as on
master, rather than while the request is deserialized.
- worker.proto pointed to the protoSegmentList query option, which this
PR removed; point it to pinot.broker.mse.enable.proto.segment.list.
- The decode benchmarks now include the first getTableSegmentsMap(), so
that the legacy decode still measures the JSON parse. Allocation per
worker is unchanged (10.88 MB legacy, 7.46 MB proto at 60k segments).
---
.../broker/broker/helix/BaseBrokerStarter.java | 3 +
.../MultiStageBrokerRequestHandler.java | 4 +-
pinot-common/src/main/proto/worker.proto | 16 ++
.../tests/MultiStageEngineIntegrationTest.java | 52 +++++
.../BaseLogicalTableIntegrationTest.java | 45 ++++
.../routing/BenchmarkSegmentListEncoding.java | 133 ++++++++++++
.../planner/physical/DispatchablePlanContext.java | 14 +-
.../planner/physical/PinotDispatchPlanner.java | 2 +
.../pinot/query/routing/QueryPlanSerDeUtils.java | 121 ++++++++++-
.../apache/pinot/query/routing/WorkerMetadata.java | 106 ++++++---
.../query/routing/QueryPlanSerDeUtilsTest.java | 239 +++++++++++++++++++++
.../query/service/dispatch/QueryDispatcher.java | 72 ++++++-
.../pinot/query/service/server/QueryServer.java | 10 +-
.../service/dispatch/QueryDispatcherTest.java | 49 +++++
.../query/service/server/QueryServerAuthzTest.java | 2 +-
.../query/service/server/QueryServerTest.java | 21 +-
.../apache/pinot/spi/utils/CommonConstants.java | 15 ++
17 files changed, 844 insertions(+), 60 deletions(-)
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/broker/helix/BaseBrokerStarter.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/broker/helix/BaseBrokerStarter.java
index a4080f8d960..f6378e38041 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/broker/helix/BaseBrokerStarter.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/broker/helix/BaseBrokerStarter.java
@@ -586,6 +586,9 @@ public abstract class BaseBrokerStarter implements
ServiceStartable {
MultiStageBrokerRequestHandler finalHandler =
multiStageBrokerRequestHandler;
_routingManager.setServerReenableCallback(
serverInstance ->
finalHandler.getQueryDispatcher().resetClientConnectionBackoff(serverInstance));
+ // Lets an operator turn the proto segment list encoding on and off
through cluster config, without a restart.
+ _clusterConfigChangeHandler.registerClusterConfigChangeListener(
+ multiStageBrokerRequestHandler.getQueryDispatcher());
}
TimeSeriesRequestHandler timeSeriesRequestHandler = null;
if
(StringUtils.isNotBlank(_brokerConf.getProperty(PinotTimeSeriesConfiguration.getEnabledLanguagesConfigKey())))
{
diff --git
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
index 299be308c2c..8c70b5cd539 100644
---
a/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
+++
b/pinot-broker/src/main/java/org/apache/pinot/broker/requesthandler/MultiStageBrokerRequestHandler.java
@@ -210,11 +210,13 @@ public class MultiStageBrokerRequestHandler extends
BaseBrokerRequestHandler {
long streamStatsDrainMs = _config.getProperty(
CommonConstants.Broker.CONFIG_OF_STREAM_STATS_DRAIN_MS,
CommonConstants.Broker.DEFAULT_STREAM_STATS_DRAIN_MS);
+ boolean enableProtoSegmentList =
_config.getProperty(CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST,
+ CommonConstants.Broker.DEFAULT_MSE_ENABLE_PROTO_SEGMENT_LIST);
_mailboxService = new MailboxService(hostname, port, InstanceType.BROKER,
config, tlsConfig);
_queryDispatcher =
new QueryDispatcher(_mailboxService, failureDetector, tlsConfig,
isQueryCancellationEnabled(), cancelTimeout,
dispatchKeepAliveTimeMs, dispatchKeepAliveTimeoutMs,
dispatchKeepAliveWithoutCalls, _streamStatsDefault,
- streamStatsDrainMs);
+ streamStatsDrainMs, enableProtoSegmentList);
LOGGER.info("Initialized MultiStageBrokerRequestHandler on host: {}, port:
{} with broker id: {}, timeout: {}ms, "
+ "query log max length: {}, query log max rate: {}, query
cancellation enabled: {}", hostname, port,
_brokerId, _brokerTimeoutMs, _queryLogger.getMaxQueryLengthToLog(),
_queryLogger.getLogRateLimit(),
diff --git a/pinot-common/src/main/proto/worker.proto
b/pinot-common/src/main/proto/worker.proto
index ceca82b65ea..2a84f200035 100644
--- a/pinot-common/src/main/proto/worker.proto
+++ b/pinot-common/src/main/proto/worker.proto
@@ -94,6 +94,22 @@ message WorkerMetadata {
int32 workedId = 1;
map<int32, bytes> mailboxInfos = 2; // Stage id to serialized MailboxInfos
map<string, string> customProperty = 3;
+ // Leaf-stage segments to scan, keyed by table type (OFFLINE / REALTIME).
Presence marks a leaf-stage worker.
+ // Sent instead of the JSON-encoded "tableSegmentsMap" custom property when
the broker enables the proto segment
+ // list encoding (see the "pinot.broker.mse.enable.proto.segment.list"
broker config); otherwise, and on older
+ // brokers, the custom property is sent.
+ SegmentsMap tableSegmentsMap = 4;
+ // Leaf-stage segments of a logical table, keyed by physical table name
(with type suffix). Same encoding rules.
+ SegmentsMap logicalTableSegmentsMap = 5;
+}
+
+// Segment names to scan per table type or physical table name, see
WorkerMetadata.
+message SegmentsMap {
+ map<string, SegmentList> segments = 1;
+}
+
+message SegmentList {
+ repeated string segment = 1;
}
message MailboxInfos {
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/MultiStageEngineIntegrationTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/MultiStageEngineIntegrationTest.java
index e5f83cb01a3..7c6d2a67417 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/MultiStageEngineIntegrationTest.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/MultiStageEngineIntegrationTest.java
@@ -48,7 +48,9 @@ import org.apache.commons.io.FileUtils;
import org.apache.commons.lang3.StringUtils;
import org.apache.helix.model.HelixConfigScope;
import org.apache.helix.model.builder.HelixConfigScopeBuilder;
+import org.apache.pinot.broker.requesthandler.BrokerRequestHandlerDelegate;
import
org.apache.pinot.controller.api.resources.PinotQueryResource.MultiStageQueryValidationRequest;
+import org.apache.pinot.query.service.dispatch.QueryDispatcher;
import org.apache.pinot.spi.config.table.HashFunction;
import org.apache.pinot.spi.config.table.RoutingConfig;
import org.apache.pinot.spi.config.table.TableConfig;
@@ -293,6 +295,56 @@ public class MultiStageEngineIntegrationTest extends
BaseClusterIntegrationTestS
List.of(), "LONG", "LONG");
}
+ /// The proto segment list encoding ships disabled, so a query uses the
legacy JSON encoding until an operator turns
+ /// it on in cluster config, which has to take effect without restarting the
broker. Both encodings must return
+ /// identical results, for queries with one and with several leaf stages.
+ @Test
+ public void testProtoSegmentListEncodingIsTransparent()
+ throws Exception {
+ QueryDispatcher dispatcher = ((BrokerRequestHandlerDelegate)
_brokerStarters.get(0).getBrokerRequestHandler())
+ .getMultiStageBrokerRequestHandler().getQueryDispatcher();
+ assertFalse(dispatcher.isEnableProtoSegmentList(), "The proto segment list
encoding must ship disabled");
+
+ String table = getTableName();
+ List<String> queries = List.of(
+ "SELECT COUNT(*) FROM " + table,
+ "SELECT Carrier, COUNT(*), MAX(ArrDelay) FROM " + table + " WHERE
DaysSinceEpoch > 16312 "
+ + "GROUP BY Carrier ORDER BY Carrier",
+ "SELECT COUNT(*) FROM " + table + " a JOIN (SELECT DISTINCT Carrier
FROM " + table + ") b "
+ + "ON a.Carrier = b.Carrier");
+
+ Map<String, JsonNode> legacyRows = new HashMap<>();
+ for (String query : queries) {
+ JsonNode response = postQuery(query);
+ assertTrue(response.get("exceptions").isEmpty(), "Unexpected exceptions
with the legacy encoding: " + response);
+ legacyRows.put(query, response.get("resultTable").get("rows"));
+ }
+
+ HelixConfigScope scope =
+ new
HelixConfigScopeBuilder(HelixConfigScope.ConfigScopeProperty.CLUSTER).forCluster(getHelixClusterName())
+ .build();
+ try {
+ // What an operator does once every server has been upgraded.
+ _helixManager.getConfigAccessor()
+ .set(scope,
CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST, "true");
+ TestUtils.waitForCondition(aVoid ->
dispatcher.isEnableProtoSegmentList(), 10_000L,
+ "Enabling the proto segment list encoding in cluster config did not
reach the broker");
+
+ for (String query : queries) {
+ JsonNode response = postQuery(query);
+ assertTrue(response.get("exceptions").isEmpty(),
+ "Unexpected exceptions with the proto encoding: " + response);
+ assertEquals(response.get("resultTable").get("rows"),
legacyRows.get(query),
+ "The segment list encoding changed the result of: " + query);
+ }
+ } finally {
+ _helixManager.getConfigAccessor()
+ .set(scope,
CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST, "false");
+ TestUtils.waitForCondition(aVoid ->
!dispatcher.isEnableProtoSegmentList(), 10_000L,
+ "Disabling the proto segment list encoding in cluster config did not
reach the broker");
+ }
+ }
+
@Test
public void testPartiallyEmptyWithPhysicalOptimizerFailsFast()
throws Exception {
diff --git
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/BaseLogicalTableIntegrationTest.java
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/BaseLogicalTableIntegrationTest.java
index f466de3a985..eeec0438444 100644
---
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/BaseLogicalTableIntegrationTest.java
+++
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/BaseLogicalTableIntegrationTest.java
@@ -29,10 +29,12 @@ import java.util.Map;
import java.util.stream.Collectors;
import java.util.stream.Stream;
import org.apache.commons.io.FileUtils;
+import org.apache.pinot.broker.requesthandler.BrokerRequestHandlerDelegate;
import org.apache.pinot.integration.tests.BaseClusterIntegrationTestSet;
import org.apache.pinot.integration.tests.ClusterIntegrationTestUtils;
import org.apache.pinot.integration.tests.QueryAssert;
import org.apache.pinot.integration.tests.QueryGenerator;
+import org.apache.pinot.query.service.dispatch.QueryDispatcher;
import org.apache.pinot.spi.config.table.QueryConfig;
import org.apache.pinot.spi.config.table.TableConfig;
import org.apache.pinot.spi.config.table.TableType;
@@ -41,6 +43,7 @@ import org.apache.pinot.spi.data.PhysicalTableConfig;
import org.apache.pinot.spi.data.Schema;
import org.apache.pinot.spi.data.TimeBoundaryConfig;
import org.apache.pinot.spi.exception.QueryErrorCode;
+import org.apache.pinot.spi.utils.CommonConstants;
import org.apache.pinot.spi.utils.JsonUtils;
import org.apache.pinot.spi.utils.builder.LogicalTableConfigBuilder;
import org.apache.pinot.spi.utils.builder.TableConfigBuilder;
@@ -480,6 +483,48 @@ public abstract class BaseLogicalTableIntegrationTest
extends BaseClusterIntegra
"Broker pruning changed the result of a logical table query");
}
+ /// Both leaf-stage segment list encodings must return the same result for a
logical table, whose leaf workers carry
+ /// `logicalTableSegmentsMap`, keyed by physical table name, rather than the
table-type keyed `tableSegmentsMap`.
+ /// The encoding is turned on through cluster config, the way an operator
would, and turned off again because the
+ /// cluster is shared with the other logical-table test classes.
+ @Test
+ public void testProtoSegmentListPreservesLogicalTableResults()
+ throws Exception {
+ setUseMultiStageQueryEngine(true);
+ String query = "SELECT Carrier, COUNT(*) FROM " + getLogicalTableName() +
" WHERE DaysSinceEpoch > 16312 "
+ + "GROUP BY Carrier ORDER BY Carrier LIMIT 100";
+ // The cluster was started by the shared suite instance, which is the one
holding the broker starter.
+ QueryDispatcher dispatcher =
+ ((BrokerRequestHandlerDelegate)
_sharedClusterTestSuite._brokerStarters.get(0).getBrokerRequestHandler())
+ .getMultiStageBrokerRequestHandler().getQueryDispatcher();
+ assertTrue(!dispatcher.isEnableProtoSegmentList(), "The proto segment list
encoding must ship disabled");
+
+ JsonNode legacy = postQuery(query);
+ assertTrue(legacy.get("exceptions").isEmpty(), "Unexpected exceptions with
the legacy encoding: " + legacy);
+
+ try {
+ setProtoSegmentList(true);
+ TestUtils.waitForCondition(aVoid ->
dispatcher.isEnableProtoSegmentList(), 10_000L,
+ "Enabling the proto segment list encoding in cluster config did not
reach the broker");
+
+ JsonNode proto = postQuery(query);
+ assertTrue(proto.get("exceptions").isEmpty(), "Unexpected exceptions
with the proto encoding: " + proto);
+ assertEquals(proto.get("resultTable").get("rows"),
legacy.get("resultTable").get("rows"),
+ "The segment list encoding changed the result of a logical table
query");
+ } finally {
+ setProtoSegmentList(false);
+ TestUtils.waitForCondition(aVoid ->
!dispatcher.isEnableProtoSegmentList(), 10_000L,
+ "Disabling the proto segment list encoding in cluster config did not
reach the broker");
+ }
+ }
+
+ private void setProtoSegmentList(boolean enabled)
+ throws Exception {
+ sendPostRequest(_controllerRequestURLBuilder.forClusterConfigs(),
+ JsonUtils.objectToString(
+
Map.of(CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST,
String.valueOf(enabled))));
+ }
+
@Test
public void testDisableGroovyQueryTableConfigOverride()
throws Exception {
diff --git
a/pinot-perf/src/main/java/org/apache/pinot/query/routing/BenchmarkSegmentListEncoding.java
b/pinot-perf/src/main/java/org/apache/pinot/query/routing/BenchmarkSegmentListEncoding.java
new file mode 100644
index 00000000000..6152f5c69cd
--- /dev/null
+++
b/pinot-perf/src/main/java/org/apache/pinot/query/routing/BenchmarkSegmentListEncoding.java
@@ -0,0 +1,133 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.query.routing;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import org.apache.pinot.common.proto.Worker;
+import org.openjdk.jmh.annotations.Benchmark;
+import org.openjdk.jmh.annotations.BenchmarkMode;
+import org.openjdk.jmh.annotations.Fork;
+import org.openjdk.jmh.annotations.Level;
+import org.openjdk.jmh.annotations.Measurement;
+import org.openjdk.jmh.annotations.Mode;
+import org.openjdk.jmh.annotations.OutputTimeUnit;
+import org.openjdk.jmh.annotations.Param;
+import org.openjdk.jmh.annotations.Scope;
+import org.openjdk.jmh.annotations.Setup;
+import org.openjdk.jmh.annotations.State;
+import org.openjdk.jmh.annotations.Threads;
+import org.openjdk.jmh.annotations.Warmup;
+import org.openjdk.jmh.runner.Runner;
+import org.openjdk.jmh.runner.RunnerException;
+import org.openjdk.jmh.runner.options.OptionsBuilder;
+
+
+/// Measures the two wire encodings of a leaf-stage worker's segment list: the
legacy JSON custom property and the
+/// native protobuf fields, on the broker (encode) and on the server (decode).
+///
+/// The measured unit is one leaf-stage worker, which is what the broker
encodes once per worker per query in
+/// [QueryPlanSerDeUtils#toProtoWorkerMetadataList] and a server decodes once
per worker in
+/// `fromProtoWorkerMetadata`. Encode includes `toByteArray`, because protobuf
defers the UTF-8 encoding of its
+/// strings to serialization while the JSON path materializes the whole string
up front; leaving it out would flatter
+/// the proto path. Decode starts from the bytes, as the server does, so it
includes the protobuf parse.
+///
+/// Lives in the `org.apache.pinot.query.routing` package so that it can call
the package-private decode entry point
+/// the server uses, rather than widening it for a benchmark.
+@BenchmarkMode(Mode.AverageTime)
+@OutputTimeUnit(TimeUnit.MICROSECONDS)
+@Warmup(iterations = 3, time = 1)
+@Measurement(iterations = 5, time = 1)
+@Fork(1)
+@Threads(1)
+@State(Scope.Benchmark)
+public class BenchmarkSegmentListEncoding {
+ /// Segment counts of one leaf-stage worker, from a small table to one
worker of a very large table.
+ @Param({"1000", "20000", "60000"})
+ private int _numSegments;
+
+ private List<WorkerMetadata> _workerMetadataList;
+ private byte[] _legacyBytes;
+ private byte[] _protoBytes;
+
+ @Setup(Level.Trial)
+ public void setUp()
+ throws Exception {
+ List<String> segments = new ArrayList<>(_numSegments);
+ for (int i = 0; i < _numSegments; i++) {
+ // A realistic 60-character offline segment name.
+ segments.add(String.format("airlineStats_OFFLINE_16071_16101_%027d", i));
+ }
+ WorkerMetadata workerMetadata = new WorkerMetadata(0,
+ Map.of(1, new MailboxInfos(new MailboxInfo("localhost", 12345,
List.of(0)))), new HashMap<>());
+ workerMetadata.setTableSegmentsMap(Map.of("OFFLINE", segments));
+ _workerMetadataList = List.of(workerMetadata);
+
+ _legacyBytes =
QueryPlanSerDeUtils.toProtoWorkerMetadataList(_workerMetadataList,
false).get(0).toByteArray();
+ _protoBytes =
QueryPlanSerDeUtils.toProtoWorkerMetadataList(_workerMetadataList,
true).get(0).toByteArray();
+ System.out.printf("%n[%d segments] wire size: legacy %d bytes, proto %d
bytes%n", _numSegments,
+ _legacyBytes.length, _protoBytes.length);
+ }
+
+ /// Build plus `toByteArray`: the whole cost the broker pays per leaf-stage
worker at dispatch.
+ @Benchmark
+ public int encodeLegacy() {
+ return QueryPlanSerDeUtils.toProtoWorkerMetadataList(_workerMetadataList,
false).get(0).toByteArray().length;
+ }
+
+ @Benchmark
+ public int encodeProto() {
+ return QueryPlanSerDeUtils.toProtoWorkerMetadataList(_workerMetadataList,
true).get(0).toByteArray().length;
+ }
+
+ /// Build only, without `toByteArray`: the part the broker used to do on the
query-compile executor at plan time.
+ @Benchmark
+ public Worker.WorkerMetadata buildLegacy() {
+ return QueryPlanSerDeUtils.toProtoWorkerMetadataList(_workerMetadataList,
false).get(0);
+ }
+
+ @Benchmark
+ public Worker.WorkerMetadata buildProto() {
+ return QueryPlanSerDeUtils.toProtoWorkerMetadataList(_workerMetadataList,
true).get(0);
+ }
+
+ /// Decode includes the first `getTableSegmentsMap()`, since the legacy JSON
is only parsed then, on the worker's own
+ /// thread; without it the legacy decode would measure nothing but the
protobuf parse.
+ @Benchmark
+ public Map<String, List<String>> decodeLegacy()
+ throws Exception {
+ return
QueryPlanSerDeUtils.fromProtoWorkerMetadata(Worker.WorkerMetadata.parseFrom(_legacyBytes))
+ .getTableSegmentsMap();
+ }
+
+ @Benchmark
+ public Map<String, List<String>> decodeProto()
+ throws Exception {
+ return
QueryPlanSerDeUtils.fromProtoWorkerMetadata(Worker.WorkerMetadata.parseFrom(_protoBytes))
+ .getTableSegmentsMap();
+ }
+
+ public static void main(String[] args)
+ throws RunnerException {
+ new Runner(new
OptionsBuilder().include(BenchmarkSegmentListEncoding.class.getSimpleName()).build()).run();
+ }
+}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/DispatchablePlanContext.java
b/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/DispatchablePlanContext.java
index 585576dfa99..58e595d131e 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/DispatchablePlanContext.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/DispatchablePlanContext.java
@@ -198,11 +198,21 @@ public class DispatchablePlanContext {
QueryServerInstance queryServerInstance = serverEntry.getValue();
serverInstanceToWorkerIdsMap.computeIfAbsent(queryServerInstance, k ->
new ArrayList<>()).add(workerId);
WorkerMetadata workerMetadata = new WorkerMetadata(workerId,
workerIdToMailboxesMap.get(workerId));
+ // A leaf-stage worker is identified by carrying a (possibly empty)
segment map, so every worker of a
+ // scanning stage has to be present in the map. Fail loudly here
instead of letting the worker decay
+ // into an intermediate-stage worker on the server. This guards an
invariant rather than an expected path:
+ // every site that populates these maps (WorkerManager and
PlanFragmentAndMailboxAssignment) fills them in
+ // the same per-worker loop as workerIdToServerInstanceMap, so every
worker gets an entry.
if (workerIdToSegmentsMap != null) {
-
workerMetadata.setTableSegmentsMap(workerIdToSegmentsMap.get(workerId));
+ Map<String, List<String>> segmentsMap =
workerIdToSegmentsMap.get(workerId);
+ Preconditions.checkNotNull(segmentsMap, "Missing segments map for
worker id: %s", workerId);
+ workerMetadata.setTableSegmentsMap(segmentsMap);
}
if (workerIdToTableNameSegmentsMap != null) {
-
workerMetadata.setLogicalTableSegmentsMap(workerIdToTableNameSegmentsMap.get(workerId));
+ Map<String, List<String>> tableNameSegmentsMap =
workerIdToTableNameSegmentsMap.get(workerId);
+ Preconditions.checkNotNull(tableNameSegmentsMap, "Missing logical
table segments map for worker id: %s",
+ workerId);
+ workerMetadata.setLogicalTableSegmentsMap(tableNameSegmentsMap);
}
workerMetadataArray[workerId] = workerMetadata;
}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/PinotDispatchPlanner.java
b/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/PinotDispatchPlanner.java
index a5baa46428b..602f817f394 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/PinotDispatchPlanner.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/query/planner/physical/PinotDispatchPlanner.java
@@ -217,6 +217,8 @@ public class PinotDispatchPlanner {
fragmentMap.put(0, reduceStage);
}
WorkerMetadata workerMetadata = workerMetadataList.get(0);
+ // Stage-0 workers never carry segment maps, and this keeps it that way
now that they live outside the custom
+ // properties: the copy below takes the custom properties only.
reduceStage.setWorkerMetadataList(List.of(
new WorkerMetadata(workerMetadata.getWorkerId(), Map.of(),
workerMetadata.getCustomProperties())));
}
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/QueryPlanSerDeUtils.java
b/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/QueryPlanSerDeUtils.java
index 955f05f8d5e..5ad7b559d85 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/QueryPlanSerDeUtils.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/QueryPlanSerDeUtils.java
@@ -18,10 +18,14 @@
*/
package org.apache.pinot.query.routing;
+import com.fasterxml.jackson.core.JsonProcessingException;
+import com.google.common.annotations.VisibleForTesting;
import com.google.common.collect.Maps;
import com.google.protobuf.ByteString;
import com.google.protobuf.InvalidProtocolBufferException;
import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
@@ -29,9 +33,21 @@ import org.apache.pinot.common.proto.Plan;
import org.apache.pinot.common.proto.Worker;
import org.apache.pinot.query.planner.plannode.PlanNode;
import org.apache.pinot.query.planner.serde.PlanNodeDeserializer;
+import org.apache.pinot.spi.utils.JsonUtils;
/// This utility class serialize/deserialize between [Worker.StagePlan]
elements to Planner elements.
+///
+/// The leaf-stage segment maps of a [WorkerMetadata] have two wire encodings,
picked per request by the broker:
+///
+/// - **proto**: the native `tableSegmentsMap` / `logicalTableSegmentsMap`
fields of [Worker.WorkerMetadata].
+/// - **legacy JSON**: a JSON string under the
[WorkerMetadata#TABLE_SEGMENTS_MAP_KEY] /
+/// [WorkerMetadata#LOGICAL_TABLE_SEGMENTS_MAP_KEY] custom property, which
is all that servers predating the proto
+/// fields understand.
+///
+/// Decoding accepts both, so a server always understands every broker; the
proto encoding is off until an operator
+/// turns it on, which they only do once every server understands it (see
+/// `CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST`).
public class QueryPlanSerDeUtils {
private QueryPlanSerDeUtils() {
}
@@ -54,15 +70,51 @@ public class QueryPlanSerDeUtils {
return new StageMetadata(protoStageMetadata.getStageId(),
workerMetadataList, customProperties);
}
- private static WorkerMetadata fromProtoWorkerMetadata(Worker.WorkerMetadata
protoWorkerMetadata)
+ @VisibleForTesting
+ static WorkerMetadata fromProtoWorkerMetadata(Worker.WorkerMetadata
protoWorkerMetadata)
throws InvalidProtocolBufferException {
Map<Integer, ByteString> protoMailboxInfosMap =
protoWorkerMetadata.getMailboxInfosMap();
Map<Integer, MailboxInfos> mailboxInfosMap =
Maps.newHashMapWithExpectedSize(protoMailboxInfosMap.size());
for (Map.Entry<Integer, ByteString> entry :
protoMailboxInfosMap.entrySet()) {
mailboxInfosMap.put(entry.getKey(),
fromProtoMailboxInfos(entry.getValue()));
}
- return new WorkerMetadata(protoWorkerMetadata.getWorkedId(),
mailboxInfosMap,
- protoWorkerMetadata.getCustomPropertyMap());
+ // A broker using the legacy encoding ships the segment maps as JSON
custom properties. Move them out of the custom
+ // properties into WorkerMetadata unparsed: each worker parses its own
list on first access, on its own thread,
+ // instead of every worker of the stage being parsed here one after
another. The custom properties stay
+ // unmodifiable either way, as the proto map view is, so that no server
path can come to depend on writing to them
+ // under one encoding only.
+ Map<String, String> customProperties =
protoWorkerMetadata.getCustomPropertyMap();
+ String tableSegmentsJson =
customProperties.get(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY);
+ String logicalTableSegmentsJson =
customProperties.get(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY);
+ if (tableSegmentsJson != null || logicalTableSegmentsJson != null) {
+ Map<String, String> strippedProperties = new HashMap<>(customProperties);
+ strippedProperties.remove(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY);
+ strippedProperties.remove(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY);
+ customProperties = Collections.unmodifiableMap(strippedProperties);
+ }
+ WorkerMetadata workerMetadata =
+ new WorkerMetadata(protoWorkerMetadata.getWorkedId(), mailboxInfosMap,
customProperties);
+ if (protoWorkerMetadata.hasTableSegmentsMap()) {
+
workerMetadata.setTableSegmentsMap(fromProtoSegmentsMap(protoWorkerMetadata.getTableSegmentsMap()));
+ } else if (tableSegmentsJson != null) {
+ workerMetadata.setTableSegmentsMapJson(tableSegmentsJson);
+ }
+ if (protoWorkerMetadata.hasLogicalTableSegmentsMap()) {
+ workerMetadata.setLogicalTableSegmentsMap(
+
fromProtoSegmentsMap(protoWorkerMetadata.getLogicalTableSegmentsMap()));
+ } else if (logicalTableSegmentsJson != null) {
+ workerMetadata.setLogicalTableSegmentsMapJson(logicalTableSegmentsJson);
+ }
+ return workerMetadata;
+ }
+
+ private static Map<String, List<String>>
fromProtoSegmentsMap(Worker.SegmentsMap protoSegmentsMap) {
+ Map<String, Worker.SegmentList> protoSegments =
protoSegmentsMap.getSegmentsMap();
+ Map<String, List<String>> segmentsMap =
Maps.newHashMapWithExpectedSize(protoSegments.size());
+ for (Map.Entry<String, Worker.SegmentList> entry :
protoSegments.entrySet()) {
+ segmentsMap.put(entry.getKey(), new
ArrayList<>(entry.getValue().getSegmentList()));
+ }
+ return segmentsMap;
}
private static MailboxInfos fromProtoMailboxInfos(ByteString
protoMailboxInfos)
@@ -76,15 +128,66 @@ public class QueryPlanSerDeUtils {
return Worker.Properties.parseFrom(protoProperties).getPropertyMap();
}
+ /// Encodes the worker metadata for the wire with the leaf-stage segment
maps in the legacy JSON encoding, which every
+ /// server understands. Kept for callers that predate the proto encoding.
public static List<Worker.WorkerMetadata>
toProtoWorkerMetadataList(List<WorkerMetadata> workerMetadataList) {
- return
workerMetadataList.stream().map(QueryPlanSerDeUtils::toProtoWorkerMetadata).collect(Collectors.toList());
+ return toProtoWorkerMetadataList(workerMetadataList, false);
+ }
+
+ /// Encodes the worker metadata for the wire, with the leaf-stage segment
maps as native proto fields when
+ /// `protoSegmentList` is set and as legacy JSON custom properties otherwise
(see the class documentation).
+ public static List<Worker.WorkerMetadata>
toProtoWorkerMetadataList(List<WorkerMetadata> workerMetadataList,
+ boolean protoSegmentList) {
+ List<Worker.WorkerMetadata> protoWorkerMetadataList = new
ArrayList<>(workerMetadataList.size());
+ for (WorkerMetadata workerMetadata : workerMetadataList) {
+ protoWorkerMetadataList.add(toProtoWorkerMetadata(workerMetadata,
protoSegmentList));
+ }
+ return protoWorkerMetadataList;
}
- private static Worker.WorkerMetadata toProtoWorkerMetadata(WorkerMetadata
workerMetadata) {
- Map<Integer, ByteString> mailboxInfosMap =
workerMetadata.getMailboxInfosMap().entrySet().stream()
- .collect(Collectors.toMap(Map.Entry::getKey, e ->
e.getValue().toProtoBytes()));
- return
Worker.WorkerMetadata.newBuilder().setWorkedId(workerMetadata.getWorkerId())
-
.putAllMailboxInfos(mailboxInfosMap).putAllCustomProperty(workerMetadata.getCustomProperties()).build();
+ private static Worker.WorkerMetadata toProtoWorkerMetadata(WorkerMetadata
workerMetadata,
+ boolean protoSegmentList) {
+ Worker.WorkerMetadata.Builder builder = Worker.WorkerMetadata.newBuilder()
+ .setWorkedId(workerMetadata.getWorkerId())
+ .putAllCustomProperty(workerMetadata.getCustomProperties());
+ for (Map.Entry<Integer, MailboxInfos> entry :
workerMetadata.getMailboxInfosMap().entrySet()) {
+ builder.putMailboxInfos(entry.getKey(), entry.getValue().toProtoBytes());
+ }
+ Map<String, List<String>> tableSegmentsMap =
workerMetadata.getTableSegmentsMap();
+ if (tableSegmentsMap != null) {
+ if (protoSegmentList) {
+ builder.setTableSegmentsMap(toProtoSegmentsMap(tableSegmentsMap));
+ } else {
+ builder.putCustomProperty(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY,
encodeSegmentsMapJson(tableSegmentsMap));
+ }
+ }
+ Map<String, List<String>> logicalTableSegmentsMap =
workerMetadata.getLogicalTableSegmentsMap();
+ if (logicalTableSegmentsMap != null) {
+ if (protoSegmentList) {
+
builder.setLogicalTableSegmentsMap(toProtoSegmentsMap(logicalTableSegmentsMap));
+ } else {
+
builder.putCustomProperty(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY,
+ encodeSegmentsMapJson(logicalTableSegmentsMap));
+ }
+ }
+ return builder.build();
+ }
+
+ private static Worker.SegmentsMap toProtoSegmentsMap(Map<String,
List<String>> segmentsMap) {
+ Worker.SegmentsMap.Builder builder = Worker.SegmentsMap.newBuilder();
+ for (Map.Entry<String, List<String>> entry : segmentsMap.entrySet()) {
+ builder.putSegments(entry.getKey(),
Worker.SegmentList.newBuilder().addAllSegment(entry.getValue()).build());
+ }
+ return builder.build();
+ }
+
+ /// JSON-encodes a segments map as `{"OFFLINE":["seg1","seg2"]}` for the
legacy encoding.
+ private static String encodeSegmentsMapJson(Map<String, List<String>>
segmentsMap) {
+ try {
+ return JsonUtils.objectToString(segmentsMap);
+ } catch (JsonProcessingException e) {
+ throw new RuntimeException("Unable to serialize segments map: " +
segmentsMap, e);
+ }
}
public static Worker.MailboxInfos toProtoMailboxInfos(List<MailboxInfo>
mailboxInfos) {
diff --git
a/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/WorkerMetadata.java
b/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/WorkerMetadata.java
index 0572a5da52d..f5781e780d5 100644
---
a/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/WorkerMetadata.java
+++
b/pinot-query-planner/src/main/java/org/apache/pinot/query/routing/WorkerMetadata.java
@@ -18,7 +18,6 @@
*/
package org.apache.pinot.query.routing;
-import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.core.type.TypeReference;
import java.io.IOException;
import java.util.HashMap;
@@ -36,20 +35,45 @@ import org.apache.pinot.spi.utils.JsonUtils;
/// - the mailbox info required to construct data transfer linkages.
/// - the partition mechanism of the data being execute on this worker.
///
+/// The segment maps are held as plain objects: they are only encoded for the
wire in [QueryPlanSerDeUtils] when a
+/// request is built for the server that runs the worker, so the planner never
pays for encoding on the compile path.
+///
+/// On a server, a segment map that arrived in the legacy JSON encoding is
kept as the raw JSON and only parsed on first
+/// access, so that each worker parses its own list on its own thread when its
leaf stage is compiled, rather than
+/// every worker of a stage being parsed one after another while the request
is deserialized.
+///
+/// Thread-safety: the raw JSON is set while the request is deserialized,
before the instance is handed to a worker.
+/// The parsed maps are published through `volatile` fields; two threads
racing on the first access may both parse the
+/// JSON, which is harmless since they produce equal maps.
+///
/// TODO: WorkerMetadata now doesn't have info directly about how to construct
the mailboxes. instead it rely on
/// MailboxSendNode and MailboxReceiveNode to derive the info during runtime.
this should changed to plan time soon.
public class WorkerMetadata {
+ /// Custom-property keys under which brokers that predate the proto segment
list encoding ship the segment maps as
+ /// JSON strings. Still written (when the proto encoding is disabled) and
read by [QueryPlanSerDeUtils] so that mixed
+ /// broker/server versions keep working; never present in
[#getCustomProperties()] of a decoded instance.
public static final String TABLE_SEGMENTS_MAP_KEY = "tableSegmentsMap";
public static final String LOGICAL_TABLE_SEGMENTS_MAP_KEY =
"logicalTableSegmentsMap";
+ private static final TypeReference<Map<String, List<String>>>
SEGMENTS_MAP_TYPE = new TypeReference<>() {
+ };
+
private final int _workerId;
private final Map<Integer, MailboxInfos> _mailboxInfosMap;
private final Map<String, String> _customProperties;
+ @Nullable
+ private volatile Map<String, List<String>> _tableSegmentsMap;
+ @Nullable
+ private volatile Map<String, List<String>> _logicalTableSegmentsMap;
+ /// The legacy JSON encoding of [#_tableSegmentsMap] as received from the
broker, parsed on first access.
+ @Nullable
+ private String _tableSegmentsMapJson;
+ /// The legacy JSON encoding of [#_logicalTableSegmentsMap] as received from
the broker, parsed on first access.
+ @Nullable
+ private String _logicalTableSegmentsMapJson;
public WorkerMetadata(int workerId, Map<Integer, MailboxInfos>
mailboxInfosMap) {
- _workerId = workerId;
- _mailboxInfosMap = mailboxInfosMap;
- _customProperties = new HashMap<>();
+ this(workerId, mailboxInfosMap, new HashMap<>());
}
public WorkerMetadata(int workerId, Map<Integer, MailboxInfos>
mailboxInfosMap,
@@ -71,52 +95,64 @@ public class WorkerMetadata {
return _customProperties;
}
+ /// Segments to scan keyed by table type (`OFFLINE` / `REALTIME`), or `null`
for a worker that scans no physical
+ /// table (intermediate stage, or a logical-table leaf).
@Nullable
public Map<String, List<String>> getTableSegmentsMap() {
- return deserializeStringSegmentListMap(TABLE_SEGMENTS_MAP_KEY);
- }
-
- private Map<String, List<String>> deserializeStringSegmentListMap(String
propertyKey) {
- String tableSegmentsMapStr = _customProperties.get(propertyKey);
- if (tableSegmentsMapStr != null) {
- try {
- return JsonUtils.stringToObject(tableSegmentsMapStr, new
TypeReference<Map<String, List<String>>>() {
- });
- } catch (IOException e) {
- throw new RuntimeException("Unable to deserialize " + propertyKey + "
: " + tableSegmentsMapStr, e);
- }
- } else {
- return null;
+ Map<String, List<String>> tableSegmentsMap = _tableSegmentsMap;
+ if (tableSegmentsMap == null && _tableSegmentsMapJson != null) {
+ tableSegmentsMap = decodeSegmentsMapJson(_tableSegmentsMapJson);
+ _tableSegmentsMap = tableSegmentsMap;
}
+ return tableSegmentsMap;
}
- public boolean isLeafStageWorker() {
- return _customProperties.containsKey(TABLE_SEGMENTS_MAP_KEY)
- || _customProperties.containsKey(LOGICAL_TABLE_SEGMENTS_MAP_KEY);
+ /// Stores `tableSegmentsMap` by reference, and it is only encoded for the
wire when the query is dispatched, so the
+ /// caller must not mutate it (or its lists) once the plan is built.
+ public void setTableSegmentsMap(Map<String, List<String>> tableSegmentsMap) {
+ _tableSegmentsMap = tableSegmentsMap;
}
- public void setTableSegmentsMap(Map<String, List<String>> tableSegmentsMap) {
- String tableSegmentsMapStr;
- try {
- tableSegmentsMapStr = JsonUtils.objectToString(tableSegmentsMap);
- } catch (JsonProcessingException e) {
- throw new RuntimeException("Unable to serialize table segments map: " +
tableSegmentsMap, e);
- }
- _customProperties.put(TABLE_SEGMENTS_MAP_KEY, tableSegmentsMapStr);
+ /// Stores the legacy JSON encoding of the table segments map, to be parsed
by [#getTableSegmentsMap] on first access.
+ void setTableSegmentsMapJson(String tableSegmentsMapJson) {
+ _tableSegmentsMapJson = tableSegmentsMapJson;
}
+ /// Segments to scan keyed by physical table name (with type suffix), or
`null` for a worker that scans no logical
+ /// table.
@Nullable
public Map<String, List<String>> getLogicalTableSegmentsMap() {
- return deserializeStringSegmentListMap(LOGICAL_TABLE_SEGMENTS_MAP_KEY);
+ Map<String, List<String>> logicalTableSegmentsMap =
_logicalTableSegmentsMap;
+ if (logicalTableSegmentsMap == null && _logicalTableSegmentsMapJson !=
null) {
+ logicalTableSegmentsMap =
decodeSegmentsMapJson(_logicalTableSegmentsMapJson);
+ _logicalTableSegmentsMap = logicalTableSegmentsMap;
+ }
+ return logicalTableSegmentsMap;
}
+ /// Stores `logicalTableSegmentsMap` by reference, with the same no-mutation
contract as [#setTableSegmentsMap].
public void setLogicalTableSegmentsMap(Map<String, List<String>>
logicalTableSegmentsMap) {
- String logicalTableSegmentsMapStr;
+ _logicalTableSegmentsMap = logicalTableSegmentsMap;
+ }
+
+ /// Stores the legacy JSON encoding of the logical table segments map, to be
parsed by
+ /// [#getLogicalTableSegmentsMap] on first access.
+ void setLogicalTableSegmentsMapJson(String logicalTableSegmentsMapJson) {
+ _logicalTableSegmentsMapJson = logicalTableSegmentsMapJson;
+ }
+
+ /// A leaf-stage worker carries a (possibly empty) segment map, parsed or
not; an intermediate-stage worker carries
+ /// none.
+ public boolean isLeafStageWorker() {
+ return _tableSegmentsMap != null || _logicalTableSegmentsMap != null ||
_tableSegmentsMapJson != null
+ || _logicalTableSegmentsMapJson != null;
+ }
+
+ private static Map<String, List<String>> decodeSegmentsMapJson(String
segmentsMapJson) {
try {
- logicalTableSegmentsMapStr =
JsonUtils.objectToString(logicalTableSegmentsMap);
- } catch (JsonProcessingException e) {
- throw new RuntimeException("Unable to serialize table segments map: " +
logicalTableSegmentsMap, e);
+ return JsonUtils.stringToObject(segmentsMapJson, SEGMENTS_MAP_TYPE);
+ } catch (IOException e) {
+ throw new RuntimeException("Unable to deserialize segments map: " +
segmentsMapJson, e);
}
- _customProperties.put(LOGICAL_TABLE_SEGMENTS_MAP_KEY,
logicalTableSegmentsMapStr);
}
}
diff --git
a/pinot-query-planner/src/test/java/org/apache/pinot/query/routing/QueryPlanSerDeUtilsTest.java
b/pinot-query-planner/src/test/java/org/apache/pinot/query/routing/QueryPlanSerDeUtilsTest.java
new file mode 100644
index 00000000000..4640a709099
--- /dev/null
+++
b/pinot-query-planner/src/test/java/org/apache/pinot/query/routing/QueryPlanSerDeUtilsTest.java
@@ -0,0 +1,239 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+package org.apache.pinot.query.routing;
+
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import javax.annotation.Nullable;
+import org.apache.pinot.common.proto.Worker;
+import org.apache.pinot.spi.utils.JsonUtils;
+import org.testng.annotations.DataProvider;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNull;
+import static org.testng.Assert.assertSame;
+import static org.testng.Assert.assertThrows;
+import static org.testng.Assert.assertTrue;
+
+
+/// Tests the two wire encodings of the leaf-stage segment maps in
[QueryPlanSerDeUtils]: the native proto fields and
+/// the legacy JSON custom properties, including a server decoding what a
pre-proto broker sends.
+public class QueryPlanSerDeUtilsTest {
+ private static final Map<String, String> CUSTOM_PROPERTIES = Map.of("foo",
"bar");
+ private static final Map<String, List<String>> TABLE_SEGMENTS_MAP =
+ Map.of("OFFLINE", List.of("seg_0", "seg_1"), "REALTIME",
List.of("seg__0__0__20240101T0000Z"));
+ private static final Map<String, List<String>> LOGICAL_TABLE_SEGMENTS_MAP =
+ Map.of("t1_OFFLINE", List.of("t1_seg_0"), "t2_REALTIME",
List.of("t2_seg_0", "t2_seg_1"));
+
+ @DataProvider
+ public static Object[][] encodings() {
+ return new Object[][]{{true}, {false}};
+ }
+
+ @Test(dataProvider = "encodings")
+ public void testLeafWorkerRoundTrip(boolean protoSegmentList)
+ throws Exception {
+ WorkerMetadata workerMetadata = leafWorker(TABLE_SEGMENTS_MAP, null);
+
+ Worker.WorkerMetadata proto = toProto(workerMetadata, protoSegmentList);
+ assertEquals(proto.hasTableSegmentsMap(), protoSegmentList);
+ assertFalse(proto.hasLogicalTableSegmentsMap());
+
assertEquals(proto.getCustomPropertyMap().containsKey(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY),
!protoSegmentList);
+
assertFalse(proto.getCustomPropertyMap().containsKey(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY));
+ assertEquals(proto.getCustomPropertyMap().get("foo"), "bar");
+
+ WorkerMetadata decoded =
QueryPlanSerDeUtils.fromProtoWorkerMetadata(proto);
+ assertEquals(decoded.getWorkerId(), 3);
+ assertEquals(decoded.getTableSegmentsMap(), TABLE_SEGMENTS_MAP);
+ assertNull(decoded.getLogicalTableSegmentsMap());
+ assertTrue(decoded.isLeafStageWorker());
+ // The JSON is never surfaced as a custom property of the decoded
metadata, whichever encoding was used.
+ assertEquals(decoded.getCustomProperties(), CUSTOM_PROPERTIES);
+ MailboxInfo mailboxInfo =
decoded.getMailboxInfosMap().get(2).getMailboxInfos().get(0);
+ assertEquals(mailboxInfo.getHostname(), "localhost");
+ assertEquals(mailboxInfo.getPort(), 1234);
+ assertEquals(mailboxInfo.getWorkerIds(), List.of(0, 1));
+ }
+
+ @Test(dataProvider = "encodings")
+ public void testLogicalTableLeafWorkerRoundTrip(boolean protoSegmentList)
+ throws Exception {
+ WorkerMetadata workerMetadata = leafWorker(null,
LOGICAL_TABLE_SEGMENTS_MAP);
+
+ Worker.WorkerMetadata proto = toProto(workerMetadata, protoSegmentList);
+ assertFalse(proto.hasTableSegmentsMap());
+ assertEquals(proto.hasLogicalTableSegmentsMap(), protoSegmentList);
+
assertEquals(proto.getCustomPropertyMap().containsKey(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY),
+ !protoSegmentList);
+
+ WorkerMetadata decoded =
QueryPlanSerDeUtils.fromProtoWorkerMetadata(proto);
+ assertNull(decoded.getTableSegmentsMap());
+ assertEquals(decoded.getLogicalTableSegmentsMap(),
LOGICAL_TABLE_SEGMENTS_MAP);
+ assertTrue(decoded.isLeafStageWorker());
+ assertEquals(decoded.getCustomProperties(), CUSTOM_PROPERTIES);
+ }
+
+ @Test(dataProvider = "encodings")
+ public void testEmptySegmentListStillMarksLeafWorker(boolean
protoSegmentList)
+ throws Exception {
+ // A padded worker of a partitioned table scans no segment but must still
run the leaf stage.
+ Map<String, List<String>> emptySegments = Map.of("OFFLINE", new
ArrayList<>());
+ WorkerMetadata decoded =
+
QueryPlanSerDeUtils.fromProtoWorkerMetadata(toProto(leafWorker(emptySegments,
null), protoSegmentList));
+ assertEquals(decoded.getTableSegmentsMap(), emptySegments);
+ assertTrue(decoded.isLeafStageWorker());
+ }
+
+ /// A segments map with no entries at all encodes, in the proto encoding, to
`SegmentsMap.getDefaultInstance()`.
+ /// The leaf/intermediate distinction rides on proto3's explicit presence
for singular message fields, which must
+ /// keep that default instance on the wire rather than drop the field, so
this goes through real bytes.
+ @Test(dataProvider = "encodings")
+ public void testZeroEntrySegmentsMapStillMarksLeafWorker(boolean
protoSegmentList)
+ throws Exception {
+ Worker.WorkerMetadata proto = toProto(leafWorker(Map.of(), Map.of()),
protoSegmentList);
+ assertEquals(proto.hasTableSegmentsMap(), protoSegmentList);
+ assertEquals(proto.hasLogicalTableSegmentsMap(), protoSegmentList);
+
+ WorkerMetadata decoded =
QueryPlanSerDeUtils.fromProtoWorkerMetadata(proto);
+ assertEquals(decoded.getTableSegmentsMap(), Map.of());
+ assertEquals(decoded.getLogicalTableSegmentsMap(), Map.of());
+ assertTrue(decoded.isLeafStageWorker());
+ }
+
+ /// Whether a server path may write to the decoded custom properties must
not depend on the encoding the broker
+ /// picked: the legacy decode strips the JSON keys from a copy, which has to
stay as unmodifiable as the proto view.
+ @Test(dataProvider = "encodings")
+ public void testDecodedCustomPropertiesAreUnmodifiable(boolean
protoSegmentList)
+ throws Exception {
+ WorkerMetadata leaf =
+
QueryPlanSerDeUtils.fromProtoWorkerMetadata(toProto(leafWorker(TABLE_SEGMENTS_MAP,
null), protoSegmentList));
+ assertThrows(UnsupportedOperationException.class, () ->
leaf.getCustomProperties().put("k", "v"));
+
+ WorkerMetadata intermediate = QueryPlanSerDeUtils.fromProtoWorkerMetadata(
+ toProto(new WorkerMetadata(1, Map.of(), new
HashMap<>(CUSTOM_PROPERTIES)), protoSegmentList));
+ assertThrows(UnsupportedOperationException.class, () ->
intermediate.getCustomProperties().put("k", "v"));
+ }
+
+ @Test(dataProvider = "encodings")
+ public void testIntermediateWorkerRoundTrip(boolean protoSegmentList)
+ throws Exception {
+ WorkerMetadata workerMetadata = new WorkerMetadata(1, Map.of(), new
HashMap<>(CUSTOM_PROPERTIES));
+
+ Worker.WorkerMetadata proto = toProto(workerMetadata, protoSegmentList);
+ assertFalse(proto.hasTableSegmentsMap());
+ assertFalse(proto.hasLogicalTableSegmentsMap());
+ assertEquals(proto.getCustomPropertyMap(), CUSTOM_PROPERTIES);
+
+ WorkerMetadata decoded =
QueryPlanSerDeUtils.fromProtoWorkerMetadata(proto);
+ assertNull(decoded.getTableSegmentsMap());
+ assertNull(decoded.getLogicalTableSegmentsMap());
+ assertFalse(decoded.isLeafStageWorker());
+ assertEquals(decoded.getCustomProperties(), CUSTOM_PROPERTIES);
+ }
+
+ /// A broker that predates the proto fields ships Jackson-encoded JSON
custom properties; a new server must decode
+ /// exactly that.
+ /// The legacy JSON is not parsed while the request is deserialized, so that
every worker parses its own list on its
+ /// own thread, as it did before the proto encoding existed. Malformed JSON
makes that observable: decoding succeeds
+ /// and the worker is still a leaf-stage worker, and only the first access
fails.
+ @Test
+ public void testLegacyJsonIsParsedOnFirstAccessOnly()
+ throws Exception {
+ Worker.WorkerMetadata proto =
Worker.WorkerMetadata.newBuilder().setWorkedId(7)
+ .putCustomProperty(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY, "{not json")
+ .putCustomProperty(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY,
"{not json either")
+ .build();
+
+ WorkerMetadata decoded =
QueryPlanSerDeUtils.fromProtoWorkerMetadata(proto);
+ assertTrue(decoded.isLeafStageWorker(), "an unparsed segment map still
marks a leaf-stage worker");
+ assertThrows(RuntimeException.class, decoded::getTableSegmentsMap);
+ assertThrows(RuntimeException.class, decoded::getLogicalTableSegmentsMap);
+ }
+
+ /// A parsed segment map is memoized, so a worker that reads it more than
once parses it once.
+ @Test
+ public void testLegacyJsonIsParsedOnce()
+ throws Exception {
+ Worker.WorkerMetadata proto =
Worker.WorkerMetadata.newBuilder().setWorkedId(7)
+ .putCustomProperty(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY,
JsonUtils.objectToString(TABLE_SEGMENTS_MAP))
+ .putCustomProperty(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY,
+ JsonUtils.objectToString(LOGICAL_TABLE_SEGMENTS_MAP))
+ .build();
+
+ WorkerMetadata decoded =
QueryPlanSerDeUtils.fromProtoWorkerMetadata(proto);
+ assertSame(decoded.getTableSegmentsMap(), decoded.getTableSegmentsMap());
+ assertSame(decoded.getLogicalTableSegmentsMap(),
decoded.getLogicalTableSegmentsMap());
+ assertEquals(decoded.getTableSegmentsMap(), TABLE_SEGMENTS_MAP);
+ assertEquals(decoded.getLogicalTableSegmentsMap(),
LOGICAL_TABLE_SEGMENTS_MAP);
+ }
+
+ @Test
+ public void testDecodesLegacyBrokerJsonCustomProperties()
+ throws Exception {
+ Worker.WorkerMetadata proto =
Worker.WorkerMetadata.newBuilder().setWorkedId(7)
+ .putCustomProperty(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY,
JsonUtils.objectToString(TABLE_SEGMENTS_MAP))
+ .putCustomProperty(WorkerMetadata.LOGICAL_TABLE_SEGMENTS_MAP_KEY,
+ JsonUtils.objectToString(LOGICAL_TABLE_SEGMENTS_MAP))
+ .putCustomProperty("foo", "bar")
+ .build();
+
+ WorkerMetadata decoded =
QueryPlanSerDeUtils.fromProtoWorkerMetadata(proto);
+ assertEquals(decoded.getWorkerId(), 7);
+ assertEquals(decoded.getTableSegmentsMap(), TABLE_SEGMENTS_MAP);
+ assertEquals(decoded.getLogicalTableSegmentsMap(),
LOGICAL_TABLE_SEGMENTS_MAP);
+ assertTrue(decoded.isLeafStageWorker());
+ assertEquals(decoded.getCustomProperties(), CUSTOM_PROPERTIES);
+ }
+
+ /// The legacy encoding must stay readable by a server that predates the
proto fields, which parses the custom
+ /// property with Jackson: pin the exact JSON shape it expects.
+ @Test
+ public void testLegacyEncodingIsTheJacksonJsonOlderServersParse()
+ throws Exception {
+ Map<String, List<String>> segmentsMap = Map.of("OFFLINE", List.of("seg_0",
"seg-1.tar.gz", "s\u00ebg_2"));
+ Worker.WorkerMetadata proto = toProto(leafWorker(segmentsMap, null),
false);
+
assertEquals(proto.getCustomPropertyMap().get(WorkerMetadata.TABLE_SEGMENTS_MAP_KEY),
+ "{\"OFFLINE\":[\"seg_0\",\"seg-1.tar.gz\",\"s\u00ebg_2\"]}");
+ }
+
+ private static WorkerMetadata leafWorker(@Nullable Map<String, List<String>>
tableSegmentsMap,
+ @Nullable Map<String, List<String>> logicalTableSegmentsMap) {
+ MailboxInfos mailboxInfos = new MailboxInfos(new MailboxInfo("localhost",
1234, List.of(0, 1)));
+ WorkerMetadata workerMetadata = new WorkerMetadata(3, Map.of(2,
mailboxInfos), new HashMap<>(CUSTOM_PROPERTIES));
+ if (tableSegmentsMap != null) {
+ workerMetadata.setTableSegmentsMap(tableSegmentsMap);
+ }
+ if (logicalTableSegmentsMap != null) {
+ workerMetadata.setLogicalTableSegmentsMap(logicalTableSegmentsMap);
+ }
+ return workerMetadata;
+ }
+
+ /// Serializes through the public list API and parses the bytes back, as the
server does.
+ private static Worker.WorkerMetadata toProto(WorkerMetadata workerMetadata,
boolean protoSegmentList)
+ throws Exception {
+ Worker.WorkerMetadata proto =
+ QueryPlanSerDeUtils.toProtoWorkerMetadataList(List.of(workerMetadata),
protoSegmentList).get(0);
+ return Worker.WorkerMetadata.parseFrom(proto.toByteString());
+ }
+}
diff --git
a/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/dispatch/QueryDispatcher.java
b/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/dispatch/QueryDispatcher.java
index 834d3b008b8..c6524bde0a0 100644
---
a/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/dispatch/QueryDispatcher.java
+++
b/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/dispatch/QueryDispatcher.java
@@ -87,6 +87,7 @@ import
org.apache.pinot.query.runtime.plan.OpChainConverterDispatcher;
import org.apache.pinot.query.runtime.plan.OpChainExecutionContext;
import org.apache.pinot.query.runtime.plan.StageStatsTreeNode;
import org.apache.pinot.query.service.dispatch.streaming.StreamingQuerySession;
+import org.apache.pinot.spi.config.provider.PinotClusterConfigChangeListener;
import org.apache.pinot.spi.exception.QueryErrorCode;
import org.apache.pinot.spi.exception.QueryException;
import org.apache.pinot.spi.query.QueryExecutionContext;
@@ -102,7 +103,9 @@ import org.slf4j.LoggerFactory;
/// `QueryDispatcher` dispatch a query to different workers.
-public class QueryDispatcher {
+public class QueryDispatcher implements PinotClusterConfigChangeListener {
+ private static final String ENABLE_PROTO_SEGMENT_LIST_KEY =
+ CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST;
private static final Logger LOGGER =
LoggerFactory.getLogger(QueryDispatcher.class);
private static final String PINOT_BROKER_QUERY_DISPATCHER_FORMAT =
"multistage-query-dispatch-%d";
/// Maximum time (ms) to wait for outstanding `OpChainComplete` stats
messages on both the success and error
@@ -133,26 +136,43 @@ public class QueryDispatcher {
/// Cluster-level default for stream-stats mode. Used as the fallback in
[#submitAndReduce] when the query
/// does not carry an explicit [QueryOptionKey#STREAM_STATS] override.
private final boolean _streamStatsDefault;
+ /// Whether leaf-stage segment lists are shipped as native protobuf fields
of the worker metadata instead of the
+ /// legacy JSON custom property. Seeded from the static broker config and
then followed live from cluster config on
+ /// [#ENABLE_PROTO_SEGMENT_LIST_KEY], so an operator can turn it on once
every server of the cluster has been
+ /// upgraded, and off again, without restarting the brokers. `volatile`
because the cluster-config callback and the
+ /// request path race; read once per query so that all servers of one query
agree.
+ private volatile boolean _enableProtoSegmentList;
public QueryDispatcher(MailboxService mailboxService, FailureDetector
failureDetector, @Nullable TlsConfig tlsConfig,
boolean enableCancellation, Duration cancelTimeout) {
this(mailboxService, failureDetector, tlsConfig, enableCancellation,
cancelTimeout,
- GrpcKeepAliveConfig.DISABLED, false,
CommonConstants.Broker.DEFAULT_STREAM_STATS_DRAIN_MS);
+ GrpcKeepAliveConfig.DISABLED, false,
CommonConstants.Broker.DEFAULT_STREAM_STATS_DRAIN_MS,
+ CommonConstants.Broker.DEFAULT_MSE_ENABLE_PROTO_SEGMENT_LIST);
}
/// Overload that accepts gRPC keep-alive settings for broker dispatch
channels. A non-positive `keepAliveTimeMs`
- /// disables keep-alive.
+ /// disables keep-alive. Kept for callers that predate the proto segment
list encoding, which they leave disabled.
public QueryDispatcher(MailboxService mailboxService, FailureDetector
failureDetector, @Nullable TlsConfig tlsConfig,
boolean enableCancellation, Duration cancelTimeout, int keepAliveTimeMs,
int keepAliveTimeoutMs,
boolean keepAliveWithoutCalls, boolean streamStatsDefault, long
statsDrainMs) {
+ this(mailboxService, failureDetector, tlsConfig, enableCancellation,
cancelTimeout, keepAliveTimeMs,
+ keepAliveTimeoutMs, keepAliveWithoutCalls, streamStatsDefault,
statsDrainMs,
+ CommonConstants.Broker.DEFAULT_MSE_ENABLE_PROTO_SEGMENT_LIST);
+ }
+
+ /// Overload that also takes the static broker config seed for the proto
segment list encoding.
+ public QueryDispatcher(MailboxService mailboxService, FailureDetector
failureDetector, @Nullable TlsConfig tlsConfig,
+ boolean enableCancellation, Duration cancelTimeout, int keepAliveTimeMs,
int keepAliveTimeoutMs,
+ boolean keepAliveWithoutCalls, boolean streamStatsDefault, long
statsDrainMs,
+ boolean enableProtoSegmentList) {
this(mailboxService, failureDetector, tlsConfig, enableCancellation,
cancelTimeout,
new GrpcKeepAliveConfig(keepAliveTimeMs, keepAliveTimeoutMs,
keepAliveWithoutCalls),
- streamStatsDefault, statsDrainMs);
+ streamStatsDefault, statsDrainMs, enableProtoSegmentList);
}
private QueryDispatcher(MailboxService mailboxService, FailureDetector
failureDetector, @Nullable TlsConfig tlsConfig,
boolean enableCancellation, Duration cancelTimeout, GrpcKeepAliveConfig
keepAliveConfig,
- boolean streamStatsDefault, long statsDrainMs) {
+ boolean streamStatsDefault, long statsDrainMs, boolean
enableProtoSegmentList) {
_cancelTimeout = cancelTimeout;
_statsDrainMs = statsDrainMs;
_mailboxService = mailboxService;
@@ -163,6 +183,7 @@ public class QueryDispatcher {
_keepAliveConfig = keepAliveConfig;
_failureDetector = failureDetector;
_streamStatsDefault = streamStatsDefault;
+ _enableProtoSegmentList = enableProtoSegmentList;
if (enableCancellation) {
_serversByQuery = new ConcurrentHashMap<>();
@@ -359,8 +380,9 @@ public class QueryDispatcher {
// that stage). The streaming observer uses this to drain the session
latch correctly when its stream errors
// before all opchains have responded.
BlockingQueue<AsyncResponse<Worker.QueryResponse>> ackQueue = new
ArrayBlockingQueue<>(serversOut.size());
+ boolean enableProtoSegmentList = _enableProtoSegmentList;
for (QueryServerInstance server : serversOut) {
- Worker.QueryRequest request = createRequest(server, stageInfos,
protoRequestMetadata);
+ Worker.QueryRequest request = createRequest(server, stageInfos,
protoRequestMetadata, enableProtoSegmentList);
int expectedForServer = 0;
for (DispatchablePlanFragment stagePlan : plansWithoutRoot) {
List<Integer> workerIds =
stagePlan.getServerInstanceToWorkerIdMap().get(server);
@@ -633,8 +655,9 @@ public class QueryDispatcher {
ByteString protoRequestMetadata =
QueryPlanSerDeUtils.toProtoProperties(requestMetadata);
// Submit the query plan to all servers in parallel
+ boolean enableProtoSegmentList = _enableProtoSegmentList;
BlockingQueue<AsyncResponse<E>> dispatchCallbacks = dispatch(sendRequest,
serverInstancesOut, deadline,
- serverInstance -> createRequest(serverInstance, stageInfos,
protoRequestMetadata));
+ serverInstance -> createRequest(serverInstance, stageInfos,
protoRequestMetadata, enableProtoSegmentList));
processResults(requestId, serverInstancesOut.size(), resultConsumer,
deadline, dispatchCallbacks);
}
@@ -698,8 +721,39 @@ public class QueryDispatcher {
}
}
+ /// Applies the proto segment list encoding set in cluster config, which
takes effect on the next query. Anything
+ /// other than `true` — the key cleared, or a value that is not a boolean —
reads as disabled, which is the legacy
+ /// encoding every server understands.
+ @Override
+ public void onChange(Set<String> changedConfigs, Map<String, String>
clusterConfigs) {
+ if (!changedConfigs.contains(ENABLE_PROTO_SEGMENT_LIST_KEY)) {
+ return;
+ }
+ String value = clusterConfigs.get(ENABLE_PROTO_SEGMENT_LIST_KEY);
+ boolean enableProtoSegmentList = value != null &&
Boolean.parseBoolean(value.trim());
+ if (enableProtoSegmentList == _enableProtoSegmentList) {
+ return;
+ }
+ _enableProtoSegmentList = enableProtoSegmentList;
+ LOGGER.info("Updated {} from: {} to: {}", ENABLE_PROTO_SEGMENT_LIST_KEY,
!enableProtoSegmentList,
+ enableProtoSegmentList);
+ if (enableProtoSegmentList) {
+ LOGGER.warn("The proto segment list encoding is now enabled. Every
server this broker dispatches to, including "
+ + "the servers of remote clusters when multi-cluster routing is
used, must already run a version that "
+ + "understands it; leaf stages routed to an older server will fail.
Set it back to false to revert.");
+ }
+ }
+
+ @VisibleForTesting
+ public boolean isEnableProtoSegmentList() {
+ return _enableProtoSegmentList;
+ }
+
+ /// Builds the request for one server: the plans of the stages it takes part
in, with only its own workers'
+ /// metadata. The leaf-stage segment lists are encoded here, once per
worker, rather than at plan time.
private static Worker.QueryRequest createRequest(QueryServerInstance
serverInstance,
- Map<DispatchablePlanFragment, StageInfo> stageInfos, ByteString
protoRequestMetadata) {
+ Map<DispatchablePlanFragment, StageInfo> stageInfos, ByteString
protoRequestMetadata,
+ boolean enableProtoSegmentList) {
Worker.QueryRequest.Builder requestBuilder =
Worker.QueryRequest.newBuilder();
requestBuilder.setVersion(PlanVersions.V1);
@@ -713,7 +767,7 @@ public class QueryDispatcher {
workerMetadataList.add(stageWorkerMetadataList.get(workerId));
}
List<Worker.WorkerMetadata> protoWorkerMetadataList =
- QueryPlanSerDeUtils.toProtoWorkerMetadataList(workerMetadataList);
+ QueryPlanSerDeUtils.toProtoWorkerMetadataList(workerMetadataList,
enableProtoSegmentList);
StageInfo stageInfo = entry.getValue();
Worker.StagePlan requestStagePlan = Worker.StagePlan.newBuilder()
diff --git
a/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/server/QueryServer.java
b/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/server/QueryServer.java
index e3c44eabc10..42dbb7344cb 100644
---
a/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/server/QueryServer.java
+++
b/pinot-query-runtime/src/main/java/org/apache/pinot/query/service/server/QueryServer.java
@@ -423,8 +423,16 @@ public class QueryServer extends
PinotQueryWorkerGrpc.PinotQueryWorkerImplBase {
}
ByteString rootAsBytes =
PlanNodeSerializer.process(explainPlan.getRootNode()).toByteString();
StageMetadata metadata = explainPlan.getStageMetadata();
+ // The segment maps are dropped rather than encoded:
QueryDispatcher#explain, the only reader of this
+ // response, reads just the root node of each stage plan, and has
since MSE explain was introduced in
+ // 1.3.0, so encoding them would cost a JSON encode of every segment
list on every EXPLAIN for nothing.
+ List<WorkerMetadata> explainWorkerMetadataList = new
ArrayList<>(metadata.getWorkerMetadataList().size());
+ for (WorkerMetadata explainWorkerMetadata :
metadata.getWorkerMetadataList()) {
+ explainWorkerMetadataList.add(new
WorkerMetadata(explainWorkerMetadata.getWorkerId(),
+ explainWorkerMetadata.getMailboxInfosMap(),
explainWorkerMetadata.getCustomProperties()));
+ }
List<Worker.WorkerMetadata> protoWorkerMetadataList =
-
QueryPlanSerDeUtils.toProtoWorkerMetadataList(metadata.getWorkerMetadataList());
+
QueryPlanSerDeUtils.toProtoWorkerMetadataList(explainWorkerMetadataList, false);
builder.addStagePlan(Worker.StagePlan.newBuilder()
.setRootNode(rootAsBytes)
.setStageMetadata(Worker.StageMetadata.newBuilder()
diff --git
a/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/dispatch/QueryDispatcherTest.java
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/dispatch/QueryDispatcherTest.java
index a557f64e26d..d0b4caae606 100644
---
a/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/dispatch/QueryDispatcherTest.java
+++
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/dispatch/QueryDispatcherTest.java
@@ -88,6 +88,55 @@ public class QueryDispatcherTest extends QueryTestSet {
Duration.ofSeconds(1));
}
+ /// The proto segment list encoding ships disabled and is turned on by an
operator through cluster config, which has
+ /// to reach the broker without a restart, and off again the same way.
+ @Test
+ public void testProtoSegmentListFollowsClusterConfig() {
+ String key =
CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST;
+ QueryDispatcher dispatcher =
+ new QueryDispatcher(Mockito.mock(MailboxService.class),
Mockito.mock(FailureDetector.class), null, false,
+ Duration.ofSeconds(1));
+ try {
+ Assert.assertFalse(dispatcher.isEnableProtoSegmentList(), "The encoding
must ship disabled");
+
+ dispatcher.onChange(Set.of(key), Map.of(key, "true"));
+ Assert.assertTrue(dispatcher.isEnableProtoSegmentList(), "Cluster config
must turn the encoding on");
+
+ dispatcher.onChange(Set.of(key), Map.of(key, "false"));
+ Assert.assertFalse(dispatcher.isEnableProtoSegmentList(), "Cluster
config must turn the encoding off again");
+
+ // A change that does not touch the key leaves it alone.
+ dispatcher.onChange(Set.of(key), Map.of(key, "TRUE"));
+ dispatcher.onChange(Set.of("some.other.key"), Map.of("some.other.key",
"x"));
+ Assert.assertTrue(dispatcher.isEnableProtoSegmentList());
+
+ // Anything that is not a boolean reads as disabled, the safe direction.
+ dispatcher.onChange(Set.of(key), Map.of(key, "SAFE"));
+ Assert.assertFalse(dispatcher.isEnableProtoSegmentList());
+ } finally {
+ dispatcher.shutdown();
+ }
+ }
+
+ /// Clearing the cluster-config key disables the encoding, whatever the
static broker config said: the fallback is
+ /// always the legacy encoding that every server understands.
+ @Test
+ public void testClearingClusterConfigDisablesTheEncoding() {
+ String key =
CommonConstants.Broker.CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST;
+ QueryDispatcher dispatcher =
+ new QueryDispatcher(Mockito.mock(MailboxService.class),
Mockito.mock(FailureDetector.class), null, false,
+ Duration.ofSeconds(1), 0, 0, false, false,
CommonConstants.Broker.DEFAULT_STREAM_STATS_DRAIN_MS, true);
+ try {
+ Assert.assertTrue(dispatcher.isEnableProtoSegmentList(), "The static
broker config seeds the value");
+
+ dispatcher.onChange(Set.of(key), Map.of());
+ Assert.assertFalse(dispatcher.isEnableProtoSegmentList(),
+ "Clearing the key must fall back to the legacy encoding");
+ } finally {
+ dispatcher.shutdown();
+ }
+ }
+
@AfterClass
public void tearDown() {
_queryDispatcher.shutdown();
diff --git
a/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerAuthzTest.java
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerAuthzTest.java
index 101e0d91a0a..3271a548f9b 100644
---
a/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerAuthzTest.java
+++
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerAuthzTest.java
@@ -155,7 +155,7 @@ public class QueryServerAuthzTest {
DispatchablePlanFragment stagePlan =
queryPlan.getQueryStageMap().get(stageId);
Plan.PlanNode rootNode =
PlanNodeSerializer.process(stagePlan.getPlanFragment().getFragmentRoot());
List<Worker.WorkerMetadata> workerMetadataList =
-
QueryPlanSerDeUtils.toProtoWorkerMetadataList(stagePlan.getWorkerMetadataList());
+
QueryPlanSerDeUtils.toProtoWorkerMetadataList(stagePlan.getWorkerMetadataList(),
false);
ByteString customProperty =
QueryPlanSerDeUtils.toProtoProperties(stagePlan.getCustomProperties());
// this particular test set requires the request to have a single
QueryServerInstance to dispatch to
diff --git
a/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerTest.java
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerTest.java
index fd8020942f6..c8c177fd51c 100644
---
a/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerTest.java
+++
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/service/server/QueryServerTest.java
@@ -235,13 +235,26 @@ public class QueryServerTest extends QueryTestSet {
@Test(dataProvider = "testSql")
public void testWorkerAcceptsWorkerRequestCorrect(String sql)
throws Exception {
+ testWorkerAcceptsWorkerRequestCorrect(sql, false);
+ }
+
+ /// Same as [#testWorkerAcceptsWorkerRequestCorrect(String)] with the
leaf-stage segment lists shipped as native
+ /// proto fields instead of the legacy JSON custom property.
+ @Test(dataProvider = "testSql")
+ public void testWorkerAcceptsProtoSegmentListRequestCorrect(String sql)
+ throws Exception {
+ testWorkerAcceptsWorkerRequestCorrect(sql, true);
+ }
+
+ private void testWorkerAcceptsWorkerRequestCorrect(String sql, boolean
protoSegmentList)
+ throws Exception {
DispatchableSubPlan queryPlan = _queryEnvironment.planQuery(sql);
Set<DispatchablePlanFragment> stagePlans =
queryPlan.getQueryStagesWithoutRoot();
// Ignore reduce stage (stage 0)
for (DispatchablePlanFragment stagePlan : stagePlans) {
int stageId = stagePlan.getPlanFragment().getFragmentId();
// only get one worker request out.
- Worker.QueryRequest queryRequest = getQueryRequest(queryPlan, stageId);
+ Worker.QueryRequest queryRequest = getQueryRequest(queryPlan, stageId,
protoSegmentList);
Map<String, String> requestMetadata =
QueryPlanSerDeUtils.fromProtoProperties(queryRequest.getMetadata());
// submit the request for testing.
@@ -321,10 +334,14 @@ public class QueryServerTest extends QueryTestSet {
}
private Worker.QueryRequest getQueryRequest(DispatchableSubPlan queryPlan,
int stageId) {
+ return getQueryRequest(queryPlan, stageId, false);
+ }
+
+ private Worker.QueryRequest getQueryRequest(DispatchableSubPlan queryPlan,
int stageId, boolean protoSegmentList) {
DispatchablePlanFragment stagePlan =
queryPlan.getQueryStageMap().get(stageId);
Plan.PlanNode rootNode =
PlanNodeSerializer.process(stagePlan.getPlanFragment().getFragmentRoot());
List<Worker.WorkerMetadata> workerMetadataList =
-
QueryPlanSerDeUtils.toProtoWorkerMetadataList(stagePlan.getWorkerMetadataList());
+
QueryPlanSerDeUtils.toProtoWorkerMetadataList(stagePlan.getWorkerMetadataList(),
protoSegmentList);
ByteString customProperty =
QueryPlanSerDeUtils.toProtoProperties(stagePlan.getCustomProperties());
// this particular test set requires the request to have a single
QueryServerInstance to dispatch to
diff --git
a/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
b/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
index a7314eee62e..1d1d36f7986 100644
--- a/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
+++ b/pinot-spi/src/main/java/org/apache/pinot/spi/utils/CommonConstants.java
@@ -668,6 +668,21 @@ public class CommonConstants {
public static final String CONFIG_OF_STREAM_STATS_DRAIN_MS =
"pinot.broker.mse.stream.stats.drain.ms";
public static final long DEFAULT_STREAM_STATS_DRAIN_MS = 50L;
+ /// Whether a multi-stage query ships its leaf-stage segment lists as
native protobuf fields of the worker
+ /// metadata, which skips a JSON encode per leaf-stage worker on the
broker and a JSON parse per worker on the
+ /// server, instead of the legacy JSON string custom property.
+ ///
+ /// Ships disabled, and must stay disabled until every server the broker
dispatches to runs a version that
+ /// understands the proto fields, including the servers of remote clusters
when multi-cluster routing is used: an
+ /// older server finds no segments under them, concludes the worker is not
a leaf-stage worker and fails the leaf
+ /// stage. Turn it on once the rolling upgrade has finished.
+ ///
+ /// Read from cluster config as well as from the static broker config,
cluster config winning and taking effect on
+ /// the next query, so it can be turned on, and off again, without
restarting the brokers. Anything in cluster
+ /// config other than `true` — the key cleared, or a value that is not a
boolean — disables it.
+ public static final String CONFIG_OF_MSE_ENABLE_PROTO_SEGMENT_LIST =
"pinot.broker.mse.enable.proto.segment.list";
+ public static final boolean DEFAULT_MSE_ENABLE_PROTO_SEGMENT_LIST = false;
+
public static final String CONFIG_OF_USE_FIXED_REPLICA =
"pinot.broker.use.fixed.replica";
public static final boolean DEFAULT_USE_FIXED_REPLICA = false;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]