This is an automated email from the ASF dual-hosted git repository. xiangfu0 pushed a commit to branch xiangfu0/integration-lanes/01-cover-subpackages in repository https://gitbox.apache.org/repos/asf/pinot.git
commit 27f7c530ea77484a861bec05ed5e8c06a4501534 Author: Xiang Fu <[email protected]> AuthorDate: Wed Sep 23 17:03:51 2026 -0700 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. --- .../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 | 175 ++++++++++------ ...alTableWithTwoRealtimeTableIntegrationTest.java | 40 ---- .../tests/suites/LogicalTableSuite.java | 46 +++++ .../pinot/query/runtime/function/NotUdf.java | 2 +- 8 files changed, 227 insertions(+), 390 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..f7085528f9c 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,65 +548,58 @@ 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"; - // Query should not throw exception - postQuery(groovyQuery); - - // Disable groovy explicitly - queryConfig = new QueryConfig(null, true, null, null, null, null); - - logicalTableConfig.setQueryConfig(queryConfig); - updateLogicalTableConfig(logicalTableConfig); + // Enable groovy for this logical table: the query must stop failing. + applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, false, null, null, null, null), + () -> { + postQuery(groovyQuery); + return true; + }, "Groovy query kept failing after groovy was enabled"); - // 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")); + // Disable groovy explicitly. + applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, true, null, null, null, null), + () -> failsWithGroovyDisabled(groovyQuery), "Groovy query kept succeeding after groovy was disabled"); - // Remove query config - logicalTableConfig.setQueryConfig(null); - updateLogicalTableConfig(logicalTableConfig); + // Removing the query config falls back to the cluster default, which also disables groovy. + 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); 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"); + // Reset to no limit. + applyQueryConfigAndAwait(logicalTableConfig, new QueryConfig(null, null, null, null, null, null), + () -> succeeds(starQuery), "Query kept failing after the response size limit was cleared"); } @Test @@ -623,33 +636,61 @@ public abstract class BaseLogicalTableIntegrationTest extends BaseClusterIntegra @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"); + // 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"); - //Reset to null. - queryConfig = new QueryConfig(null, null, null, null, null, null); + // 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"); + + // Reset to no timeout override. + 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 +703,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 --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
