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]

Reply via email to