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]