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 36a3253a2d6 Run the integration test subpackages that no CI lane ever 
selected (#19652)
36a3253a2d6 is described below

commit 36a3253a2d683c67ba10a189da322ec8d88f407a
Author: Xiang Fu <[email protected]>
AuthorDate: Thu Sep 24 11:59:03 2026 -0700

    Run the integration test subpackages that no CI lane ever selected (#19652)
    
    * Run the integration test subpackages that no CI lane ever selected
    
    The lane include patterns are single-level globs such as
    "org/apache/pinot/integration/tests/L*Test.java", and a "*" never crosses a
    directory boundary, so classes in subpackages were never matched. Only tpch/
    and server/realtime/ (explicit entries), custom/ (CustomClusterSuite) and 
the
    two Kinesis classes (KinesisSuite) ever ran. logicaltable/, multicluster/,
    legacy/, udf/ and realtime/ingestion/ did not, and
    CancelQueryIntegrationTests was missed as well because its name ends in
    "Tests", which no A-Z prefix matches.
    
    The subpackages are added to the lanes, balanced against the last healthy 
lane
    runtimes (set-1 lane-a 2524s, lane-b 2555s, set-2 lane-a 2055s, lane-b 
2231s):
    
      set-1 lane-b  legacy/
      set-2 lane-a  LogicalTableSuite, KafkaPartitionSubsetChaosIntegrationTest
      set-2 lane-b  multicluster/
    
    logicaltable/ needs a focused suite execution rather than a lane include:
    BaseLogicalTableIntegrationTest starts its cluster from @BeforeSuite, so its
    subclasses share that cluster only when they run as one suite. 
LogicalTableSuite
    selects the whole package, the way CustomClusterSuite does for custom/, so a
    class added later cannot be silently skipped. It excludes
    KafkaPartitionSubsetChaosIntegrationTest, which lives in the package but 
does
    not extend that base and starts its own cluster from @BeforeClass; inside 
the
    suite that cluster would sit alongside the shared one for its whole 
duration,
    and it timed out there while passing on its own. A lane selects it directly.
    
    Enabling these classes exposed five failures, all fixed here:
    
    - NotUdf looked up LogicalFunctions.not(boolean), but #17189 changed the
      signature to not(Boolean). The constructor threw NoSuchMethodException, so
      every ServiceLoader.load(Udf.class) failed with ServiceConfigurationError 
and
      UdfTest could not run at all. Only the UDF test framework loads that SPI, 
so
      no query path was affected.
    - LogicalTableWithTwoRealtimeTableIntegrationTest could not create its 
tables.
      getKafkaTopic() defaults to the class simple name, so every subclass has 
its
      own topic, but @BeforeSuite runs on one instance and creates only that
      instance's topic. This class's topic was left to be auto-created with a 
single
      partition, so the controller rejected its second table for pinning
      stream.kafka.partition.ids=1. Each class now creates its own topic with 
its
      own partition count before its table configs are validated.
    - The suite failed in teardown even with every test passing. tearDown() 
purges
      the shared cluster through cleanup(), which reads Helix and the property 
store
      directly, but only the @BeforeSuite instance held those handles. The 
resulting
      NullPointerException also replaced the message cleanup() was about to 
report.
      Non-owner subclasses now inherit those handles.
    - testLogicalTableWithEmptyOfflineTable still expected numServersQueried == 
1
      for the multi-stage engine. #18538 made the multi-stage broker 
short-circuit
      when all leaf stages are empty, so both engines now report 0.
    - testQueryTimeOut, testMaxQueryResponseSizeTableConfig and
      testDisableGroovyQueryTableConfigOverride each applied a query config 
through
      the controller and then asserted on the very next query.
      updateLogicalTableConfig is asynchronous - brokers observe the change 
through
      a ZooKeeper property store listener - so the assertion raced propagation.
      They now wait for the new behavior through a shared helper. 
testQueryTimeOut
      also accepts every stage's timeout code, which lets
      LogicalTableWithTwoRealtimeTableIntegrationTest drop its own copy of the 
test
      instead of maintaining a list that omitted EXECUTION_TIMEOUT.
    
    Verified by running LogicalTableSuite: 138 tests, no failures.
    
    Three classes that had also never been selected stay out, each with the 
reason
    recorded next to the pattern, because they fail for their own pre-existing
    reasons rather than anything in this change:
    
    - UdfTest: its snapshots under src/test/resources/udf-test-results went 
stale
      while the class was unrunnable. Most of the drift is additive, but 
refreshing
      them would also record that transform initialization errors now surface 
as a
      generic "Operator execution error" instead of naming the cause, which 
wants a
      deliberate decision first.
    - CancelQueryIntegrationTests: testCancelByClientQueryId observes no
      QueryCancellationError on either engine.
    - realtime/ingestion/KafkaIncreaseDecreasePartitionsIntegrationTest: POST
      /tables hangs until the 60s client timeout when the test adds its second
      realtime table. Reproduced twice.
    
    Finally, .pinot_tests_integration.sh, .pinot_tests_custom_integration.sh and
    .pinot_tests_kinesis_integration.sh are referenced by no workflow, and the
    kinesis one gates on RUN_TEST_SET==2 although KinesisSuite runs in set-1
    lane-a. They are removed together with the integration-tests-set-1 and
    integration-tests-set-2 profiles, whose only consumer they were and whose
    duplicated include lists are what let the lanes drift unnoticed. The four 
lane
    profiles are now the single source of truth.
    
    * Address review: cover the NotUdf fix and stop the config waits passing 
vacuously
    
    Three findings from the review:
    
    - NotUdf's fix was invisible to CI, because UdfTest is the only selected 
test
      that loads the Udf SPI and it stays quarantined. UdfServiceLoaderTest now
      iterates the SPI, which is the failure mode a stale reflective lookup
      actually produces: one broken provider makes ServiceLoader throw
      ServiceConfigurationError for every caller, not just for that provider.
      Iterating also covers implementations added later. A second test pins why
      the boxed lookup is required: LogicalFunctions.not(Boolean) exists and the
      primitive not(boolean) the stale lookup asked for does not. Reintroducing
      the old lookup fails both tests, with ServiceConfigurationError and
      NoSuchMethodException respectively.
    
    - testMaxServerResponseSizeTableConfig was still updating the logical table
      config and querying immediately at each of its three steps. It races 
broker
      propagation exactly like the three tests already converted, so it now uses
      applyQueryConfigAndAwait too.
    
    - The reset step of each of those tests waited for an outcome that already
      held under the config it was replacing, so the wait returned immediately
      and proved nothing. Each test now restores the restrictive config before
      clearing it, which makes every transition flip the observable outcome.
      Groovy is checked from the enabled state because the cluster default
      disables it exactly like the explicit override does.
    
    The query strings also stop hardcoding "mytable" now that these tests run 
for
    every subclass.
    
    Verified by running LogicalTableSuite: 146 tests, no failures.
---
 .../pr-tests/.pinot_tests_custom_integration.sh    |  33 ---
 .../scripts/pr-tests/.pinot_tests_integration.sh   |  67 ------
 .../pr-tests/.pinot_tests_kinesis_integration.sh   |  33 ---
 pinot-integration-tests/pom.xml                    | 221 +++++++-------------
 .../BaseLogicalTableIntegrationTest.java           | 228 +++++++++++++--------
 ...alTableWithTwoRealtimeTableIntegrationTest.java |  40 ----
 .../tests/suites/LogicalTableSuite.java            |  46 +++++
 .../pinot/query/runtime/function/NotUdf.java       |   2 +-
 .../runtime/function/UdfServiceLoaderTest.java     |  70 +++++++
 9 files changed, 325 insertions(+), 415 deletions(-)

diff --git 
a/.github/workflows/scripts/pr-tests/.pinot_tests_custom_integration.sh 
b/.github/workflows/scripts/pr-tests/.pinot_tests_custom_integration.sh
deleted file mode 100755
index 0de5d5fdc74..00000000000
--- a/.github/workflows/scripts/pr-tests/.pinot_tests_custom_integration.sh
+++ /dev/null
@@ -1,33 +0,0 @@
-#!/bin/bash -x
-#
-# 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.
-#
-
-# Java version
-java -version
-
-# Check network
-ifconfig
-netstat -i
-
-# Custom Integration Tests
-cd pinot-integration-tests || exit 1
-if [ "$RUN_TEST_SET" == "1" ]; then
-  mvn test \
-      -P github-actions,codecoverage,custom-cluster-integration-test-suite || 
exit 1
-fi
diff --git a/.github/workflows/scripts/pr-tests/.pinot_tests_integration.sh 
b/.github/workflows/scripts/pr-tests/.pinot_tests_integration.sh
deleted file mode 100755
index ca3ea68e63f..00000000000
--- a/.github/workflows/scripts/pr-tests/.pinot_tests_integration.sh
+++ /dev/null
@@ -1,67 +0,0 @@
-#!/bin/bash -x
-#
-# 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.
-#
-
-# Java version
-java -version
-
-# Check network
-ifconfig
-netstat -i
-
-print_surefire_dumps() {
-  local reports_dir="target/surefire-reports"
-  local dump_files
-
-  if [ ! -d "${reports_dir}" ]; then
-    echo "Surefire reports directory not found: ${reports_dir}"
-    return
-  fi
-
-  dump_files="$(find "${reports_dir}" -maxdepth 1 -type f \
-    \( -name "*.dump" -o -name "*.dumpstream" -o -name "*jvmRun*" \) | sort)"
-  if [ -z "${dump_files}" ]; then
-    echo "No Surefire dump files found under ${reports_dir}"
-    return
-  fi
-
-  echo "Printing Surefire dump files from ${reports_dir}"
-  while IFS= read -r dump_file; do
-    [ -z "${dump_file}" ] && continue
-    echo "===== BEGIN ${dump_file} ====="
-    cat "${dump_file}"
-    echo "===== END ${dump_file} ====="
-  done <<< "${dump_files}"
-}
-
-# Integration Tests
-cd pinot-integration-tests || exit 1
-case "$RUN_TEST_SET" in
-  1|2)
-    mvn test \
-        -P "github-actions,codecoverage,integration-tests-set-${RUN_TEST_SET}" 
|| {
-      print_surefire_dumps
-      exit 1
-    }
-    ;;
-  *)
-    echo "Unsupported RUN_TEST_SET value: ${RUN_TEST_SET}"
-    exit 1
-    ;;
-esac
diff --git 
a/.github/workflows/scripts/pr-tests/.pinot_tests_kinesis_integration.sh 
b/.github/workflows/scripts/pr-tests/.pinot_tests_kinesis_integration.sh
deleted file mode 100755
index 518724c3247..00000000000
--- a/.github/workflows/scripts/pr-tests/.pinot_tests_kinesis_integration.sh
+++ /dev/null
@@ -1,33 +0,0 @@
-#!/bin/bash -x
-#
-# 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.
-#
-
-# Java version
-java -version
-
-# Check network
-ifconfig
-netstat -i
-
-# Kinesis ingestion integration tests
-cd pinot-integration-tests || exit 1
-if [ "$RUN_TEST_SET" == "2" ]; then
-  mvn test \
-      -P github-actions,codecoverage,kinesis-integration-test-suite || exit 1
-fi
diff --git a/pinot-integration-tests/pom.xml b/pinot-integration-tests/pom.xml
index c5b0c79e610..18a6d351fe5 100644
--- a/pinot-integration-tests/pom.xml
+++ b/pinot-integration-tests/pom.xml
@@ -100,160 +100,14 @@
         
<pinot.integration.test.skip.named.suites>true</pinot.integration.test.skip.named.suites>
       </properties>
     </profile>
-    <!--
-      Keep every top-level A-Z test prefix in exactly one set. 
TenantRebalanceIntegrationTest runs in set 1 to
-      balance the observed GitHub Actions runtime across the two jobs. The 
CI-only lane profiles below own the
-      focused data-source suites and their finer-grained balancing.
-    -->
-    <profile>
-      <id>integration-tests-set-1</id>
-      <activation>
-        <activeByDefault>false</activeByDefault>
-      </activation>
-      <properties>
-        <argLine>${pinot.integration.test.jvm.args}</argLine>
-      </properties>
-      <build>
-        <plugins>
-          <plugin>
-            <groupId>org.apache.maven.plugins</groupId>
-            <artifactId>maven-surefire-plugin</artifactId>
-            <configuration>
-              <skipTests>false</skipTests>
-              <includes>
-                
<include>org/apache/pinot/integration/tests/A*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/B*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/C*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/D*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/E*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/F*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/G*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/H*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/I*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/J*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/K*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/L*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/M*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/N*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/TenantRebalanceIntegrationTest.java</include>
-              </includes>
-              <excludes>
-                
<exclude>org/apache/pinot/integration/tests/DateTimeFieldSpecHybridClusterIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/ExactlyOnceKafkaRealtimeClusterIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/HybridClusterIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/IngestionConfigHybridIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/KafkaConfluentSchemaRegistryAvroMessageDecoderRealtimeClusterIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/KafkaConsumingSegmentToBeMovedSummaryIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/LLCRealtimeClusterIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/MultiNodesOfflineClusterIntegrationTest.java</exclude>
-              </excludes>
-            </configuration>
-            <executions>
-              <execution>
-                <id>shared-kafka-realtime-integration-test-suite</id>
-                <phase>test</phase>
-                <goals>
-                  <goal>test</goal>
-                </goals>
-                <configuration>
-                  
<skipTests>${pinot.integration.test.skip.named.suites}</skipTests>
-                  <reportNameSuffix>shared-kafka</reportNameSuffix>
-                  <excludes combine.self="override">
-                    <exclude>none</exclude>
-                  </excludes>
-                  <includes combine.self="override">
-                    
<include>org/apache/pinot/integration/tests/suites/SharedKafkaRealtimeSuite.java</include>
-                  </includes>
-                </configuration>
-              </execution>
-              <execution>
-                <id>multi-nodes-offline-integration-test-suite</id>
-                <phase>test</phase>
-                <goals>
-                  <goal>test</goal>
-                </goals>
-                <configuration>
-                  
<skipTests>${pinot.integration.test.skip.named.suites}</skipTests>
-                  <reportNameSuffix>multi-nodes-offline</reportNameSuffix>
-                  <excludes combine.self="override">
-                    <exclude>none</exclude>
-                  </excludes>
-                  <includes combine.self="override">
-                    
<include>org/apache/pinot/integration/tests/suites/MultiNodesOfflineSuite.java</include>
-                  </includes>
-                </configuration>
-              </execution>
-              <execution>
-                <id>shared-hybrid-integration-test-suite</id>
-                <phase>test</phase>
-                <goals>
-                  <goal>test</goal>
-                </goals>
-                <configuration>
-                  
<skipTests>${pinot.integration.test.skip.named.suites}</skipTests>
-                  <reportNameSuffix>shared-hybrid</reportNameSuffix>
-                  <excludes combine.self="override">
-                    <exclude>none</exclude>
-                  </excludes>
-                  <includes combine.self="override">
-                    
<include>org/apache/pinot/integration/tests/suites/SharedHybridSuite.java</include>
-                  </includes>
-                </configuration>
-              </execution>
-            </executions>
-          </plugin>
-        </plugins>
-      </build>
-    </profile>
-    <profile>
-      <id>integration-tests-set-2</id>
-      <activation>
-        <activeByDefault>false</activeByDefault>
-      </activation>
-      <properties>
-        <argLine>${pinot.integration.test.jvm.args}</argLine>
-      </properties>
-      <build>
-        <plugins>
-          <plugin>
-            <groupId>org.apache.maven.plugins</groupId>
-            <artifactId>maven-surefire-plugin</artifactId>
-            <configuration>
-              <skipTests>false</skipTests>
-              <includes>
-                
<include>org/apache/pinot/integration/tests/O*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/P*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/Q*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/R*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/S*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/T*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/U*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/V*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/W*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/X*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/Y*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/Z*Test.java</include>
-                
<include>org/apache/pinot/integration/tests/tpch/*Test.java</include>
-                <include>org/apache/pinot/server/realtime/**</include>
-              </includes>
-              <excludes>
-                
<exclude>org/apache/pinot/integration/tests/PinotLLCRealtimeSegmentManagerIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/QueryThreadContextIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/SpoolIntegrationTest.java</exclude>
-                
<exclude>org/apache/pinot/integration/tests/TenantRebalanceIntegrationTest.java</exclude>
-              </excludes>
-            </configuration>
-          </plugin>
-        </plugins>
-      </build>
-    </profile>
     <!--
       The CI-only lane profiles split each integration-test set into two 
isolated Maven processes. Each source test
       is selected by exactly one default execution or by one focused suite 
execution.
 
-      Surefire include/exclude patterns are relative to the test classes 
directory and must NOT be prefixed with
-      "**/": under the JUnit Platform provider a leading "**/" no longer 
matches zero directories, so a pattern
-      like "**/org/apache/pinot/..." silently selects nothing.
+      Patterns are relative to the test classes directory and must NOT be 
prefixed with "**/": under the JUnit
+      Platform provider a leading "**/" no longer matches zero directories, so 
"**/org/apache/pinot/..." silently
+      selects nothing. A "*" never crosses a directory boundary, so each 
subpackage needs its own entry; keep the
+      A-Z prefixes for classes directly under 
org/apache/pinot/integration/tests and list subpackages explicitly.
     -->
     <profile>
       <id>integration-tests-set-1-lane-a</id>
@@ -279,12 +133,30 @@
                 
<include>org/apache/pinot/integration/tests/F*Test.java</include>
                 
<include>org/apache/pinot/integration/tests/G*Test.java</include>
                 
<include>org/apache/pinot/integration/tests/H*Test.java</include>
+                
<include>org/apache/pinot/integration/tests/udf/*Test.java</include>
+                <!--
+                  Never selected by any lane before, and each fails for a 
pre-existing reason unrelated to the
+                  selection fix, so they stay out until their own failure is 
resolved:
+                    - CancelQueryIntegrationTests (also missed because it ends 
in "Tests", which no A-Z prefix
+                      matches): testCancelByClientQueryId sees no 
QueryCancellationError on either engine.
+                    - 
realtime/ingestion/KafkaIncreaseDecreasePartitionsIntegrationTest: POST /tables 
hangs until
+                      the 60s client timeout when the test adds its second 
realtime table.
+                  The rest of realtime/ingestion is the Kinesis pair, which 
this lane already runs through
+                  KinesisSuite below; selecting that package here would run 
them twice.
+                -->
               </includes>
               <excludes>
                 
<exclude>org/apache/pinot/integration/tests/DateTimeFieldSpecHybridClusterIntegrationTest.java</exclude>
                 
<exclude>org/apache/pinot/integration/tests/ExactlyOnceKafkaRealtimeClusterIntegrationTest.java</exclude>
                 
<exclude>org/apache/pinot/integration/tests/GrpcBrokerClusterIntegrationTest.java</exclude>
                 
<exclude>org/apache/pinot/integration/tests/HybridClusterIntegrationTest.java</exclude>
+                <!--
+                  UdfTest is quarantined, not unselected: its snapshots under 
src/test/resources/udf-test-results
+                  went stale while the class was unrunnable, and refreshing 
them would also record that transform
+                  initialization errors now surface as a generic "Operator 
execution error" instead of naming the
+                  cause. That is a maintainer's call, so the class waits for 
it; anything else added under udf/ runs.
+                -->
+                
<exclude>org/apache/pinot/integration/tests/udf/UdfTest.java</exclude>
               </excludes>
             </configuration>
             <executions>
@@ -351,6 +223,7 @@
                 
<include>org/apache/pinot/integration/tests/N*Test.java</include>
                 
<include>org/apache/pinot/integration/tests/GrpcBrokerClusterIntegrationTest.java</include>
                 
<include>org/apache/pinot/integration/tests/TenantRebalanceIntegrationTest.java</include>
+                
<include>org/apache/pinot/integration/tests/legacy/*Test.java</include>
               </includes>
               <excludes>
                 
<exclude>org/apache/pinot/integration/tests/IngestionConfigHybridIntegrationTest.java</exclude>
@@ -437,6 +310,11 @@
                 
<include>org/apache/pinot/integration/tests/P*Test.java</include>
                 
<include>org/apache/pinot/integration/tests/Q*Test.java</include>
                 
<include>org/apache/pinot/integration/tests/R*Test.java</include>
+                <!--
+                  In the logicaltable package but not a 
BaseLogicalTableIntegrationTest subclass: it runs its own
+                  cluster, so LogicalTableSuite below excludes it and it runs 
here instead.
+                -->
+                
<include>org/apache/pinot/integration/tests/logicaltable/KafkaPartitionSubsetChaosIntegrationTest.java</include>
               </includes>
               <excludes>
                 
<exclude>org/apache/pinot/integration/tests/OfflineClusterIntegrationTest.java</exclude>
@@ -445,6 +323,25 @@
                 
<exclude>org/apache/pinot/integration/tests/RealtimeConsumptionRateLimiterClusterIntegrationTest.java</exclude>
               </excludes>
             </configuration>
+            <executions>
+              <execution>
+                <id>logical-table-integration-test-suite</id>
+                <phase>test</phase>
+                <goals>
+                  <goal>test</goal>
+                </goals>
+                <configuration>
+                  
<skipTests>${pinot.integration.test.skip.named.suites}</skipTests>
+                  <reportNameSuffix>logical-table</reportNameSuffix>
+                  <excludes combine.self="override">
+                    <exclude>none</exclude>
+                  </excludes>
+                  <includes combine.self="override">
+                    
<include>org/apache/pinot/integration/tests/suites/LogicalTableSuite.java</include>
+                  </includes>
+                </configuration>
+              </execution>
+            </executions>
           </plugin>
         </plugins>
       </build>
@@ -476,6 +373,7 @@
                 
<include>org/apache/pinot/integration/tests/OfflineClusterIntegrationTest.java</include>
                 
<include>org/apache/pinot/integration/tests/RealtimeConsumptionRateLimiterClusterIntegrationTest.java</include>
                 
<include>org/apache/pinot/integration/tests/tpch/*Test.java</include>
+                
<include>org/apache/pinot/integration/tests/multicluster/*Test.java</include>
                 <include>org/apache/pinot/server/realtime/**</include>
               </includes>
               <excludes>
@@ -510,6 +408,29 @@
         </plugins>
       </build>
     </profile>
+    <profile>
+      <id>logical-table-integration-test-suite</id>
+      <activation>
+        <activeByDefault>false</activeByDefault>
+      </activation>
+      <properties>
+        <argLine>${pinot.integration.test.jvm.args}</argLine>
+      </properties>
+      <build>
+        <plugins>
+          <plugin>
+            <groupId>org.apache.maven.plugins</groupId>
+            <artifactId>maven-surefire-plugin</artifactId>
+            <configuration>
+              <skipTests>false</skipTests>
+              <includes combine.self="override">
+                
<include>org/apache/pinot/integration/tests/suites/LogicalTableSuite.java</include>
+              </includes>
+            </configuration>
+          </plugin>
+        </plugins>
+      </build>
+    </profile>
     <profile>
       <id>kinesis-integration-test-suite</id>
       <activation>
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 eeec0438444..07c84d222cb 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
@@ -20,6 +20,7 @@ package org.apache.pinot.integration.tests.logicaltable;
 
 import com.fasterxml.jackson.databind.JsonNode;
 import java.io.File;
+import java.time.Duration;
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.HashSet;
@@ -28,6 +29,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.stream.Collectors;
 import java.util.stream.Stream;
+import javax.annotation.Nullable;
 import org.apache.commons.io.FileUtils;
 import org.apache.pinot.broker.requesthandler.BrokerRequestHandlerDelegate;
 import org.apache.pinot.integration.tests.BaseClusterIntegrationTestSet;
@@ -61,7 +63,6 @@ import org.testng.annotations.Test;
 
 import static org.testng.Assert.assertEquals;
 import static org.testng.Assert.assertTrue;
-import static org.testng.Assert.expectThrows;
 
 
 public abstract class BaseLogicalTableIntegrationTest extends 
BaseClusterIntegrationTestSet {
@@ -70,6 +71,9 @@ public abstract class BaseLogicalTableIntegrationTest extends 
BaseClusterIntegra
   private static final String DEFAULT_LOGICAL_TABLE_NAME = "mytable";
   protected static final String DEFAULT_TABLE_NAME = "physicalTable";
   protected static final String EMPTY_OFFLINE_TABLE_NAME = "empty_o";
+  private static final String GROOVY_DISABLED_MESSAGE = "Groovy transform 
functions are disabled for queries";
+  private static final long CONFIG_PROPAGATION_CHECK_INTERVAL_MS = 100L;
+  private static final long CONFIG_PROPAGATION_TIMEOUT_MS = 60_000L;
   protected static BaseLogicalTableIntegrationTest _sharedClusterTestSuite = 
null;
   protected List<File> _avroFiles;
 
@@ -115,6 +119,14 @@ public abstract class BaseLogicalTableIntegrationTest 
extends BaseClusterIntegra
       _helixResourceManager = _sharedClusterTestSuite._helixResourceManager;
       _kafkaStarters = _sharedClusterTestSuite._kafkaStarters;
       _controllerBaseApiUrl = _sharedClusterTestSuite._controllerBaseApiUrl;
+      // tearDown() purges the shared cluster through cleanup() so the next 
class starts from an empty cluster,
+      // and cleanup() reads Helix and the property store directly. Only the 
instance that ran @BeforeSuite has
+      // those handles, so without copying them here cleanup() throws a 
NullPointerException that replaces
+      // whatever it was about to report.
+      _helixManager = _sharedClusterTestSuite._helixManager;
+      _helixDataAccessor = _sharedClusterTestSuite._helixDataAccessor;
+      _helixAdmin = _sharedClusterTestSuite._helixAdmin;
+      _propertyStore = _sharedClusterTestSuite._propertyStore;
     }
 
     _avroFiles = getAllAvroFiles();
@@ -142,6 +154,14 @@ public abstract class BaseLogicalTableIntegrationTest 
extends BaseClusterIntegra
 
     // create realtime table
     Map<String, List<File>> realtimeTableDataFiles = 
getRealtimeTableDataFiles();
+    if (!realtimeTableDataFiles.isEmpty()) {
+      // getKafkaTopic() defaults to the class simple name, so every subclass 
has its own topic, but the shared
+      // cluster's @BeforeSuite runs on a single instance and therefore only 
creates that one instance's topic.
+      // Create this class's topic explicitly - a no-op when it already exists 
- so that a table config pinning
+      // stream.kafka.partition.ids is validated against the expected 
partition count rather than against a
+      // single-partition topic auto-created by the broker on first access.
+      createKafkaTopic(getKafkaTopic(), getNumKafkaPartitions());
+    }
     for (Map.Entry<String, List<File>> entry : 
realtimeTableDataFiles.entrySet()) {
       String tableName = entry.getKey();
       List<File> avroFilesForTable = entry.getValue();
@@ -528,128 +548,152 @@ public abstract class BaseLogicalTableIntegrationTest 
extends BaseClusterIntegra
   @Test
   public void testDisableGroovyQueryTableConfigOverride()
       throws Exception {
-    QueryConfig queryConfig = new QueryConfig(null, false, null, null, null, 
null);
     LogicalTableConfig logicalTableConfig = 
getLogicalTableConfig(getLogicalTableName());
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-
     String groovyQuery = "SELECT 
GROOVY('{\"returnType\":\"STRING\",\"isSingleValue\":true}', "
-        + "'arg0 + arg1', FlightNum, Origin) FROM mytable";
+        + "'arg0 + arg1', FlightNum, Origin) FROM " + getLogicalTableName();
 
-    // Query should not throw exception
-    postQuery(groovyQuery);
+    // Every step has to flip the outcome of the step before it. A wait whose 
condition already holds under the
+    // previous config returns immediately and proves nothing, so the config 
being cleared is checked from the
+    // enabled state: the cluster default disables Groovy exactly like the 
explicit override does.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, false, 
null, null, null, null),
+        () -> succeeds(groovyQuery), "Groovy query kept failing after groovy 
was enabled");
 
-    // Disable groovy explicitly
-    queryConfig = new QueryConfig(null, true, null, null, null, null);
-
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, true, 
null, null, null, null),
+        () -> failsWithGroovyDisabled(groovyQuery), "Groovy query kept 
succeeding after groovy was disabled");
 
-    // grpc and http throw different exceptions. So only check error message.
-    Exception athrows = expectThrows(Exception.class, () -> 
postQuery(groovyQuery));
-    assertTrue(athrows.getMessage().contains("Groovy transform functions are 
disabled for queries"));
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, false, 
null, null, null, null),
+        () -> succeeds(groovyQuery), "Groovy query kept failing after groovy 
was enabled again");
 
-    // Remove query config
-    logicalTableConfig.setQueryConfig(null);
-    updateLogicalTableConfig(logicalTableConfig);
+    // Removing the query config falls back to the cluster default, which 
disables groovy again.
+    applyQueryConfigAndAwait(logicalTableConfig, null, () -> 
failsWithGroovyDisabled(groovyQuery),
+        "Groovy query kept succeeding after the query config was removed");
+  }
 
-    athrows = expectThrows(Exception.class, () -> postQuery(groovyQuery));
-    assertTrue(athrows.getMessage().contains("Groovy transform functions are 
disabled for queries"));
+  /// Returns whether `query` fails because groovy is disabled, rather than 
succeeding or failing for another reason.
+  ///
+  /// Returns instead of asserting so that [#applyQueryConfigAndAwait] can 
retry: an assertion failure is an
+  /// `AssertionError`, which the wait helper does not treat as a 
not-yet-satisfied condition.
+  private boolean failsWithGroovyDisabled(String query) {
+    try {
+      postQuery(query);
+      return false;
+    } catch (Exception e) {
+      // grpc and http throw different exceptions, so only check the error 
message.
+      String message = e.getMessage();
+      return message != null && message.contains(GROOVY_DISABLED_MESSAGE);
+    }
   }
 
   @Test
   public void testMaxQueryResponseSizeTableConfig()
       throws Exception {
-    String starQuery = "SELECT * from mytable";
-
-    QueryConfig queryConfig = new QueryConfig(null, null, null, null, 100L, 
null);
+    String starQuery = "SELECT * from " + getLogicalTableName();
     LogicalTableConfig logicalTableConfig = 
getLogicalTableConfig(getLogicalTableName());
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
 
-    JsonNode response = postQuery(starQuery);
-    JsonNode exceptions = response.get("exceptions");
-    assertTrue(!exceptions.isEmpty()
-        && exceptions.get(0).get("errorCode").asInt() == 
QueryErrorCode.QUERY_CANCELLATION.getId());
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, 100L, null),
+        () -> failsWith(starQuery, QueryErrorCode.QUERY_CANCELLATION),
+        "Query was not cancelled under a 100 byte response size limit");
 
-    // Query Succeeds with a high limit.
-    queryConfig = new QueryConfig(null, null, null, null, 1000000L, null);
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
+    // Query succeeds with a high limit.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, 1000000L, null),
+        () -> succeeds(starQuery), "Query kept failing under a high response 
size limit");
 
-    //Reset to null.
-    queryConfig = new QueryConfig(null, null, null, null, null, null);
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
+    // Restore the restrictive limit so that clearing it below is observable: 
waiting for success straight after
+    // the high limit would be satisfied by the high limit itself.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, 100L, null),
+        () -> failsWith(starQuery, QueryErrorCode.QUERY_CANCELLATION),
+        "Query was not cancelled after the 100 byte response size limit was 
restored");
+
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, null, null),
+        () -> succeeds(starQuery), "Query kept failing after the response size 
limit was cleared");
   }
 
   @Test
   public void testMaxServerResponseSizeTableConfig()
       throws Exception {
-    String starQuery = "SELECT * from mytable";
-
-    QueryConfig queryConfig = new QueryConfig(null, null, null, null, null, 
1000L);
+    String starQuery = "SELECT * from " + getLogicalTableName();
     LogicalTableConfig logicalTableConfig = 
getLogicalTableConfig(getLogicalTableName());
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    JsonNode response = postQuery(starQuery);
-    JsonNode exceptions = response.get("exceptions");
-    assertTrue(!exceptions.isEmpty()
-        && exceptions.get(0).get("errorCode").asInt() == 
QueryErrorCode.QUERY_CANCELLATION.getId());
 
-    // Query Succeeds with a high limit.
-    queryConfig = new QueryConfig(null, null, null, null, null, 1000000L);
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, null, 1000L),
+        () -> failsWith(starQuery, QueryErrorCode.QUERY_CANCELLATION),
+        "Query was not cancelled under a 1000 byte server response size 
limit");
 
-    //Reset to null.
-    queryConfig = new QueryConfig(null, null, null, null, null, null);
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
+    // Query succeeds with a high limit.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, null, 1000000L),
+        () -> succeeds(starQuery), "Query kept failing under a high server 
response size limit");
+
+    // Restore the restrictive limit so that clearing it below is observable.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, null, 1000L),
+        () -> failsWith(starQuery, QueryErrorCode.QUERY_CANCELLATION),
+        "Query was not cancelled after the 1000 byte server response size 
limit was restored");
+
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, null, null),
+        () -> succeeds(starQuery), "Query kept failing after the server 
response size limit was cleared");
   }
 
   @Test
   public void testQueryTimeOut()
       throws Exception {
-    String starQuery = "SELECT * from mytable";
-    QueryConfig queryConfig = new QueryConfig(1L, null, null, null, null, 
null);
+    String starQuery = "SELECT * from " + getLogicalTableName();
     LogicalTableConfig logicalTableConfig = 
getLogicalTableConfig(getLogicalTableName());
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    JsonNode response = postQuery(starQuery);
-    JsonNode exceptions = response.get("exceptions");
-    assertTrue(
-        !exceptions.isEmpty() && (exceptions.get(0).get("errorCode").asInt() 
== QueryErrorCode.BROKER_TIMEOUT.getId()
-            // Timeout may occur just before submitting the request. Then this 
error code is thrown.
-            || exceptions.get(0).get("errorCode").asInt() == 
QueryErrorCode.SERVER_NOT_RESPONDING.getId()));
-
-    // Query Succeeds with a high limit.
-    queryConfig = new QueryConfig(1000000L, null, null, null, null, null);
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
 
-    //Reset to null.
-    queryConfig = new QueryConfig(null, null, null, null, null, null);
+    // A 1 ms budget can expire at any stage, and each stage reports its own 
code: before the request is
+    // submitted (SERVER_NOT_RESPONDING), while it waits to be scheduled 
(QUERY_SCHEDULING_TIMEOUT), while a
+    // server runs it (EXECUTION_TIMEOUT) or while the broker waits for 
servers (BROKER_TIMEOUT). Which one wins
+    // depends on how much work the table shape implies, so accept any of them.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(1L, null, 
null, null, null, null),
+        () -> failsWith(starQuery, QueryErrorCode.BROKER_TIMEOUT, 
QueryErrorCode.SERVER_NOT_RESPONDING,
+            QueryErrorCode.QUERY_SCHEDULING_TIMEOUT, 
QueryErrorCode.EXECUTION_TIMEOUT),
+        "Query did not time out under a 1 ms timeout");
+
+    // Query succeeds with a high timeout.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(1000000L, 
null, null, null, null, null),
+        () -> succeeds(starQuery), "Query kept failing under a high timeout");
+
+    // Restore the 1 ms timeout so that clearing the override below is 
observable.
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(1L, null, 
null, null, null, null),
+        () -> failsWith(starQuery, QueryErrorCode.BROKER_TIMEOUT, 
QueryErrorCode.SERVER_NOT_RESPONDING,
+            QueryErrorCode.QUERY_SCHEDULING_TIMEOUT, 
QueryErrorCode.EXECUTION_TIMEOUT),
+        "Query did not time out after the 1 ms timeout was restored");
+
+    applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, 
null, null, null, null),
+        () -> succeeds(starQuery), "Query kept failing after the timeout 
override was cleared");
+  }
+
+  /// Applies `queryConfig` to `logicalTableConfig` and waits until the 
brokers act on it.
+  ///
+  /// Updating a logical table config through the controller is asynchronous: 
brokers observe the change through
+  /// a ZooKeeper property store listener, so a query issued right after the 
REST call can still be planned with
+  /// the previous config. Waiting for the new behavior keeps these assertions 
from racing that propagation.
+  private void applyQueryConfigAndAwait(LogicalTableConfig logicalTableConfig, 
@Nullable QueryConfig queryConfig,
+      TestUtils.SupplierWithException<Boolean> brokerAppliedConfig, String 
message)
+      throws Exception {
     logicalTableConfig.setQueryConfig(queryConfig);
     updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
+    TestUtils.waitForCondition(brokerAppliedConfig, 
CONFIG_PROPAGATION_CHECK_INTERVAL_MS,
+        CONFIG_PROPAGATION_TIMEOUT_MS, message, 
Duration.ofMillis(CONFIG_PROPAGATION_TIMEOUT_MS / 4));
+  }
+
+  /// Returns whether `query` completes without any exception.
+  private boolean succeeds(String query)
+      throws Exception {
+    return postQuery(query).get("exceptions").isEmpty();
+  }
+
+  /// Returns whether `query` reports an exception whose error code is one of 
`expectedErrorCodes`.
+  private boolean failsWith(String query, QueryErrorCode... expectedErrorCodes)
+      throws Exception {
+    JsonNode exceptions = postQuery(query).get("exceptions");
+    if (exceptions.isEmpty()) {
+      return false;
+    }
+    int errorCode = exceptions.get(0).get("errorCode").asInt();
+    for (QueryErrorCode expectedErrorCode : expectedErrorCodes) {
+      if (errorCode == expectedErrorCode.getId()) {
+        return true;
+      }
+    }
+    return false;
   }
 
   @Test(dataProvider = "useBothQueryEngines")
@@ -662,7 +706,9 @@ public abstract class BaseLogicalTableIntegrationTest 
extends BaseClusterIntegra
     // Query should return empty result
     JsonNode queryResponse = postQuery("SELECT count(*) FROM " + 
logicalTableName);
     assertEquals(queryResponse.get("numDocsScanned").asInt(), 0);
-    assertEquals(queryResponse.get("numServersQueried").asInt(), 
useMultiStageQueryEngine ? 1 : 0);
+    // Neither engine dispatches to servers for an empty table: the 
multi-stage broker short-circuits when all leaf
+    // stages are empty (#18538).
+    assertEquals(queryResponse.get("numServersQueried").asInt(), 0, "Query 
should not dispatch to servers");
     assertTrue(queryResponse.get("exceptions").isEmpty());
   }
 
diff --git 
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/LogicalTableWithTwoRealtimeTableIntegrationTest.java
 
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/LogicalTableWithTwoRealtimeTableIntegrationTest.java
index f09ebdd35e4..265cc2e657d 100644
--- 
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/LogicalTableWithTwoRealtimeTableIntegrationTest.java
+++ 
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/logicaltable/LogicalTableWithTwoRealtimeTableIntegrationTest.java
@@ -18,7 +18,6 @@
  */
 package org.apache.pinot.integration.tests.logicaltable;
 
-import com.fasterxml.jackson.databind.JsonNode;
 import com.google.common.primitives.Longs;
 import java.io.ByteArrayOutputStream;
 import java.io.File;
@@ -37,11 +36,9 @@ import org.apache.kafka.clients.producer.KafkaProducer;
 import org.apache.kafka.clients.producer.ProducerConfig;
 import org.apache.kafka.clients.producer.ProducerRecord;
 import org.apache.pinot.plugin.inputformat.avro.AvroUtils;
-import org.apache.pinot.spi.config.table.QueryConfig;
 import org.apache.pinot.spi.config.table.TableConfig;
 import org.apache.pinot.spi.config.table.ingestion.IngestionConfig;
 import org.apache.pinot.spi.config.table.ingestion.StreamIngestionConfig;
-import org.apache.pinot.spi.exception.QueryErrorCode;
 import org.apache.pinot.spi.utils.builder.TableNameBuilder;
 import org.apache.pinot.util.TestUtils;
 import org.testng.Assert;
@@ -49,7 +46,6 @@ import org.testng.annotations.BeforeClass;
 import org.testng.annotations.Test;
 
 import static org.testng.Assert.assertEquals;
-import static org.testng.Assert.assertTrue;
 
 
 public class LogicalTableWithTwoRealtimeTableIntegrationTest extends 
BaseLogicalTableIntegrationTest {
@@ -128,42 +124,6 @@ public class 
LogicalTableWithTwoRealtimeTableIntegrationTest extends BaseLogical
     assertEquals(_table0RecordCount + _table1RecordCount, 
getCurrentCountStarResult(LOGICAL_TABLE_NAME));
   }
 
-  @Override
-  @Test
-  public void testQueryTimeOut()
-      throws Exception {
-    String starQuery = "SELECT * from " + getLogicalTableName();
-    QueryConfig queryConfig = new QueryConfig(1L, null, null, null, null, 
null);
-    var logicalTableConfig = getLogicalTableConfig(getLogicalTableName());
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    JsonNode response = postQuery(starQuery);
-    JsonNode exceptions = response.get("exceptions");
-    if (!exceptions.isEmpty()) {
-      int errorCode = exceptions.get(0).get("errorCode").asInt();
-      assertTrue(errorCode == QueryErrorCode.BROKER_TIMEOUT.getId()
-          || errorCode == QueryErrorCode.SERVER_NOT_RESPONDING.getId()
-          || errorCode == QueryErrorCode.QUERY_SCHEDULING_TIMEOUT.getId(),
-          "Unexpected error code: " + errorCode);
-    }
-
-    // Query succeeds with a high limit.
-    queryConfig = new QueryConfig(1000000L, null, null, null, null, null);
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
-
-    // Reset to null.
-    queryConfig = new QueryConfig(null, null, null, null, null, null);
-    logicalTableConfig.setQueryConfig(queryConfig);
-    updateLogicalTableConfig(logicalTableConfig);
-    response = postQuery(starQuery);
-    exceptions = response.get("exceptions");
-    assertTrue(exceptions.isEmpty(), "Query should not throw exception");
-  }
-
   @Override
   protected TableConfig createRealtimeTableConfig(File sampleAvroFile) {
     TableConfig tableConfig = super.createRealtimeTableConfig(sampleAvroFile);
diff --git 
a/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/suites/LogicalTableSuite.java
 
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/suites/LogicalTableSuite.java
new file mode 100644
index 00000000000..537aadd0e03
--- /dev/null
+++ 
b/pinot-integration-tests/src/test/java/org/apache/pinot/integration/tests/suites/LogicalTableSuite.java
@@ -0,0 +1,46 @@
+/**
+ * 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.integration.tests.suites;
+
+import org.junit.platform.suite.api.ExcludeClassNamePatterns;
+import org.junit.platform.suite.api.IncludeClassNamePatterns;
+import org.junit.platform.suite.api.IncludeEngines;
+import org.junit.platform.suite.api.SelectPackages;
+import org.junit.platform.suite.api.Suite;
+
+
+/// Runs the logical-table scenarios with one shared cluster.
+///
+/// `BaseLogicalTableIntegrationTest` starts its cluster from `@BeforeSuite`, 
so its subclasses only share that
+/// cluster when they run in a single suite. Selecting the whole package keeps 
every class in one fork and keeps a
+/// newly added class from being silently skipped.
+///
+/// `KafkaPartitionSubsetChaosIntegrationTest` is excluded because it does not 
extend
+/// `BaseLogicalTableIntegrationTest`: it starts and stops its own cluster 
from `@BeforeClass`, which inside this
+/// suite would run alongside the shared cluster for the whole of its (long) 
duration. The lane profile selects it
+/// directly instead, so it gets a fork where only its own cluster is up.
+///
+/// Stateless suite definition; the selected tests run sequentially in one 
fork.
+@Suite
+@IncludeEngines("testng")
+@SelectPackages("org.apache.pinot.integration.tests.logicaltable")
+@IncludeClassNamePatterns(".*")
+@ExcludeClassNamePatterns(".*\\.KafkaPartitionSubsetChaosIntegrationTest")
+public class LogicalTableSuite {
+}
diff --git 
a/pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/function/NotUdf.java
 
b/pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/function/NotUdf.java
index 754d910e093..9a2699e98fb 100644
--- 
a/pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/function/NotUdf.java
+++ 
b/pinot-query-runtime/src/main/java/org/apache/pinot/query/runtime/function/NotUdf.java
@@ -40,7 +40,7 @@ public class NotUdf extends Udf.FromAnnotatedMethod {
 
   public NotUdf()
       throws NoSuchMethodException {
-    super(LogicalFunctions.class.getMethod("not", boolean.class));
+    super(LogicalFunctions.class.getMethod("not", Boolean.class));
   }
 
   @Override
diff --git 
a/pinot-query-runtime/src/test/java/org/apache/pinot/query/runtime/function/UdfServiceLoaderTest.java
 
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/runtime/function/UdfServiceLoaderTest.java
new file mode 100644
index 00000000000..d97bb34af08
--- /dev/null
+++ 
b/pinot-query-runtime/src/test/java/org/apache/pinot/query/runtime/function/UdfServiceLoaderTest.java
@@ -0,0 +1,70 @@
+/**
+ * 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.runtime.function;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.ServiceLoader;
+import org.apache.pinot.common.function.scalar.LogicalFunctions;
+import org.apache.pinot.core.udf.Udf;
+import org.testng.annotations.Test;
+
+import static org.testng.Assert.assertEquals;
+import static org.testng.Assert.assertFalse;
+import static org.testng.Assert.assertNotNull;
+import static org.testng.Assert.assertThrows;
+import static org.testng.Assert.assertTrue;
+
+
+/// Guards the [Udf] service registrations against a scalar function signature 
changing underneath them.
+///
+/// A [Udf.FromAnnotatedMethod] subclass binds to its scalar function 
reflectively in its constructor, so changing
+/// that method's signature makes the constructor throw. `ServiceLoader` turns 
that into a
+/// `ServiceConfigurationError`, which fails every caller that iterates the 
SPI rather than only the provider at
+/// fault, so one stale lookup disables UDF loading altogether.
+public class UdfServiceLoaderTest {
+
+  /// Iterating the SPI must instantiate every registered implementation. This 
is the failure mode a broken
+  /// reflective lookup produces, and it covers implementations added later 
without listing them here.
+  @Test
+  public void testEveryRegisteredUdfCanBeInstantiated() {
+    List<String> udfNames = new ArrayList<>();
+    for (Udf udf : ServiceLoader.load(Udf.class)) {
+      assertNotNull(udf.getMainName(), udf.getClass().getName() + " has no 
main name");
+      udfNames.add(udf.getClass().getName());
+    }
+    assertFalse(udfNames.isEmpty(), "No Udf implementation was registered");
+    assertTrue(udfNames.contains(NotUdf.class.getName()), "NotUdf is not 
registered: " + udfNames);
+  }
+
+  /// `LogicalFunctions.not` takes a nullable `Boolean`, and no primitive 
`boolean` overload exists, so `NotUdf`
+  /// has to look up the boxed one. Constructing it is the assertion: the 
constructor resolves that method
+  /// reflectively and throws `NoSuchMethodException` when the lookup names a 
signature that is not there.
+  @Test
+  public void testNotUdfBindsToTheNullableBooleanOverload()
+      throws NoSuchMethodException {
+    assertNotNull(LogicalFunctions.class.getMethod("not", Boolean.class));
+    // Pin the reason the boxed lookup is required: the primitive overload the 
stale lookup asked for is gone.
+    assertThrows(NoSuchMethodException.class, () -> 
LogicalFunctions.class.getMethod("not", boolean.class));
+
+    NotUdf notUdf = new NotUdf();
+    assertEquals(notUdf.getMainName(), "not");
+    assertTrue(notUdf.getAllNames().contains("not"), "NotUdf must expose the 
'not' name: " + notUdf.getAllNames());
+  }
+}


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to