github-actions[bot] commented on code in PR #68775: URL: https://github.com/apache/doris/pull/68775#discussion_r4220410219
########## be/test/service/lance_index_prewarm_test.cpp: ########## @@ -0,0 +1,123 @@ +// 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. + +#include <gtest/gtest.h> + +#include <chrono> +#include <future> +#include <utility> + +#include "common/config.h" +#include "load/channel/load_stream_mgr.h" +#include "runtime/exec_env.h" +#include "service/internal_service.h" + +namespace doris { +namespace { + +class PrewarmClosure : public google::protobuf::Closure { +public: + void Run() override { completed.set_value(); } + std::promise<void> completed; +}; + +class LanceIndexPrewarmServiceTest : public testing::Test { +protected: + void SetUp() override { + auto heavy = std::exchange(config::brpc_heavy_work_pool_threads, 1); + auto light = std::exchange(config::brpc_light_work_pool_threads, 1); + auto peer = std::exchange(config::brpc_peer_fetch_pool_threads, 1); + auto flight = std::exchange(config::brpc_arrow_flight_work_pool_threads, 1); + auto* env = ExecEnv::GetInstance(); + previous_load_mgr = std::move(env->_load_stream_mgr); + env->_load_stream_mgr = std::make_unique<LoadStreamMgr>(1); Review Comment: [P1] Make the new BE service test compile. This fixture directly moves and assigns private `ExecEnv::_load_stream_mgr` here and later accesses protected `PInternalService` pools, even though it only derives from `testing::Test`. The file is globbed into `doris_be_test`, so that target fails compilation before these tests can run. Use public test hooks or suitable friend/test subclass access for both objects. ########## regression-test/suites/external_table_p0/lance/test_lance_index_prewarm.groovy: ########## @@ -0,0 +1,143 @@ +// 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. + +import groovy.json.JsonSlurper + +import org.apache.doris.regression.util.DebugPoint +import org.apache.doris.regression.util.NodeType + +suite("test_lance_index_prewarm", "p0,external") { + if (!"true".equalsIgnoreCase(context.config.otherConfigs.get("enableIcebergTest"))) { + logger.info("Lance prewarm requires the Iceberg MinIO fixtures") + return + } + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_lance_index_prewarm" + String tableName = "${catalogName}.doris.vs_ivf_flat_f32" + sql """DROP CATALOG IF EXISTS `${catalogName}`""" + sql """ + CREATE CATALOG `${catalogName}` PROPERTIES ( + "type" = "lance", + "lance.catalog.type" = "filesystem", + "warehouse" = "s3://warehouse/lance", + "s3.endpoint" = "http://${externalEnvIp}:${minioPort}", + "s3.access_key" = "admin", + "s3.secret_key" = "password", + "s3.region" = "us-east-1", + "use_path_style" = "true" + ) + """ + sql "SET enable_file_scanner_v2 = true" + def entries = sql """SELECT IndexName, DatasetVersion FROM lance_index_entries("table"="${tableName}")""" + assertFalse(entries.isEmpty()) + String indexName = entries[0][0] + String statement = "WARM UP INDEX `${indexName}` ON ${tableName}" + String query = """ + SELECT row_id, _distance FROM vector_search( + "table"="${tableName}", "column"="embedding", + "query_vector"="[0,1,2,3,4,5,6,7,8,9,10,11,12,13,14,15]", + "top_k"="10", "metric"="l2", "nprobes"="4", "use_index"="true") + ORDER BY _distance, row_id + """ + def before = sql(query) + def warm = sql(statement) + assertEquals(1, warm.size()) + assertEquals(indexName, warm[0][1]) + // Physical entries record the index build version; prewarm pins the current dataset snapshot. + assertTrue(warm[0][2].toLong() >= entries[0][1].toLong()) + assertTrue(warm[0][3].toInteger() > 0) + assertEquals(before, sql(query), "Prewarm must preserve query results") + def retry = sql(statement) + assertEquals(warm[0][2], retry[0][2], "Prewarm must not change the dataset version") + assertEquals(warm[0][3], retry[0][3], "Retry must cover the same eligible backend set") + + // Binary PREPARE and EXECUTE must agree on the result schema, including repeated executions. + // Force server preparation even when the JDBC URL defaults to client-side emulation. + def prepared = context.getConnection().unwrap(com.mysql.cj.jdbc.JdbcConnection) + .serverPrepareStatement(statement) + try { + assertTrue(prepared instanceof com.mysql.cj.jdbc.ServerPreparedStatement) + assertEquals(5, prepared.getMetaData().getColumnCount()) + assertEquals("DatasetVersion", prepared.getMetaData().getColumnName(3)) + for (int execution = 0; execution < 2; execution++) { + def result = exec(prepared) + assertEquals(1, result.size()) + assertEquals(indexName, result[0][1]) + assertEquals(warm[0][2].toString(), result[0][2].toString()) + assertEquals(warm[0][3].toString(), result[0][3].toString()) + } + } finally { + prepared.close() + } + + test { + sql "WARM UP INDEX missing_index ON ${tableName}" + exception "does not exist in dataset version" + } + if (isCloudMode()) { + def groups = sql "SHOW CLUSTERS" + assertFalse(groups.isEmpty()) + def explicit = sql "${statement} WITH COMPUTE GROUP `${groups[0][0]}`" Review Comment: [P2] Choose an available compute group for this success assertion. `SHOW CLUSTERS` is sorted by name and includes groups with zero or unavailable backends, so `groups[0][0]` can differ from the healthy current group used by the preceding successful prewarm. In that valid multi-group state, `selectBackends()` throws `No available backends` here and the assertion fails. Select the current group from `is_current`, or choose a group with an eligible BE. ########## regression-test/suites/external_table_p0/lance/test_lance_index_prewarm.groovy: ########## @@ -0,0 +1,143 @@ +// 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. + +import groovy.json.JsonSlurper + +import org.apache.doris.regression.util.DebugPoint +import org.apache.doris.regression.util.NodeType + +suite("test_lance_index_prewarm", "p0,external") { + if (!"true".equalsIgnoreCase(context.config.otherConfigs.get("enableIcebergTest"))) { + logger.info("Lance prewarm requires the Iceberg MinIO fixtures") + return + } + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_lance_index_prewarm" + String tableName = "${catalogName}.doris.vs_ivf_flat_f32" + sql """DROP CATALOG IF EXISTS `${catalogName}`""" + sql """ + CREATE CATALOG `${catalogName}` PROPERTIES ( + "type" = "lance", + "lance.catalog.type" = "filesystem", + "warehouse" = "s3://warehouse/lance", + "s3.endpoint" = "http://${externalEnvIp}:${minioPort}", + "s3.access_key" = "admin", + "s3.secret_key" = "password", + "s3.region" = "us-east-1", + "use_path_style" = "true" + ) + """ + sql "SET enable_file_scanner_v2 = true" + def entries = sql """SELECT IndexName, DatasetVersion FROM lance_index_entries("table"="${tableName}")""" + assertFalse(entries.isEmpty()) + String indexName = entries[0][0] + String statement = "WARM UP INDEX `${indexName}` ON ${tableName}" + String query = """ + SELECT row_id, _distance FROM vector_search( + "table"="${tableName}", "column"="embedding", + "query_vector"="[0,1,2,3,4,5,6,7,8,9,10,11,12,13,14,15]", + "top_k"="10", "metric"="l2", "nprobes"="4", "use_index"="true") + ORDER BY _distance, row_id + """ + def before = sql(query) + def warm = sql(statement) + assertEquals(1, warm.size()) + assertEquals(indexName, warm[0][1]) + // Physical entries record the index build version; prewarm pins the current dataset snapshot. + assertTrue(warm[0][2].toLong() >= entries[0][1].toLong()) + assertTrue(warm[0][3].toInteger() > 0) + assertEquals(before, sql(query), "Prewarm must preserve query results") + def retry = sql(statement) + assertEquals(warm[0][2], retry[0][2], "Prewarm must not change the dataset version") + assertEquals(warm[0][3], retry[0][3], "Retry must cover the same eligible backend set") + + // Binary PREPARE and EXECUTE must agree on the result schema, including repeated executions. + // Force server preparation even when the JDBC URL defaults to client-side emulation. + def prepared = context.getConnection().unwrap(com.mysql.cj.jdbc.JdbcConnection) + .serverPrepareStatement(statement) + try { + assertTrue(prepared instanceof com.mysql.cj.jdbc.ServerPreparedStatement) + assertEquals(5, prepared.getMetaData().getColumnCount()) + assertEquals("DatasetVersion", prepared.getMetaData().getColumnName(3)) + for (int execution = 0; execution < 2; execution++) { + def result = exec(prepared) + assertEquals(1, result.size()) + assertEquals(indexName, result[0][1]) + assertEquals(warm[0][2].toString(), result[0][2].toString()) + assertEquals(warm[0][3].toString(), result[0][3].toString()) + } + } finally { + prepared.close() + } + + test { + sql "WARM UP INDEX missing_index ON ${tableName}" + exception "does not exist in dataset version" + } + if (isCloudMode()) { + def groups = sql "SHOW CLUSTERS" + assertFalse(groups.isEmpty()) + def explicit = sql "${statement} WITH COMPUTE GROUP `${groups[0][0]}`" + assertTrue(explicit[0][3].toInteger() > 0) + } + + // Fail one target only: successful prewarm on other BEs must not hide the failure. + def backends = sql_return_maparray "SHOW BACKENDS" + def target = backends.find { Review Comment: [P2] Choose the injected backend from the prewarm-eligible set. This filter accepts an alive BE with `SystemDecommissioned=false` even while it is decommissioning, but `needLoadAvailable()` excludes decommissioning BEs. If `find` picks that BE, the command can warm another BE and succeed without hitting the debug point, so the failure assertion fails independently of the cloud tag-key fix. Use the same eligibility policy as `LanceIndexPrewarm.selectBackends()`. ########## regression-test/suites/external_table_p0/lance/test_lance_index_prewarm.groovy: ########## @@ -0,0 +1,143 @@ +// 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. + +import groovy.json.JsonSlurper + +import org.apache.doris.regression.util.DebugPoint +import org.apache.doris.regression.util.NodeType + +suite("test_lance_index_prewarm", "p0,external") { + if (!"true".equalsIgnoreCase(context.config.otherConfigs.get("enableIcebergTest"))) { + logger.info("Lance prewarm requires the Iceberg MinIO fixtures") + return + } + String externalEnvIp = context.config.otherConfigs.get("externalEnvIp") + String minioPort = context.config.otherConfigs.get("iceberg_minio_port") + String catalogName = "test_lance_index_prewarm" + String tableName = "${catalogName}.doris.vs_ivf_flat_f32" + sql """DROP CATALOG IF EXISTS `${catalogName}`""" + sql """ + CREATE CATALOG `${catalogName}` PROPERTIES ( + "type" = "lance", + "lance.catalog.type" = "filesystem", + "warehouse" = "s3://warehouse/lance", + "s3.endpoint" = "http://${externalEnvIp}:${minioPort}", + "s3.access_key" = "admin", + "s3.secret_key" = "password", + "s3.region" = "us-east-1", + "use_path_style" = "true" + ) + """ + sql "SET enable_file_scanner_v2 = true" + def entries = sql """SELECT IndexName, DatasetVersion FROM lance_index_entries("table"="${tableName}")""" + assertFalse(entries.isEmpty()) + String indexName = entries[0][0] + String statement = "WARM UP INDEX `${indexName}` ON ${tableName}" + String query = """ + SELECT row_id, _distance FROM vector_search( + "table"="${tableName}", "column"="embedding", + "query_vector"="[0,1,2,3,4,5,6,7,8,9,10,11,12,13,14,15]", + "top_k"="10", "metric"="l2", "nprobes"="4", "use_index"="true") + ORDER BY _distance, row_id + """ + def before = sql(query) + def warm = sql(statement) + assertEquals(1, warm.size()) + assertEquals(indexName, warm[0][1]) + // Physical entries record the index build version; prewarm pins the current dataset snapshot. + assertTrue(warm[0][2].toLong() >= entries[0][1].toLong()) + assertTrue(warm[0][3].toInteger() > 0) + assertEquals(before, sql(query), "Prewarm must preserve query results") + def retry = sql(statement) + assertEquals(warm[0][2], retry[0][2], "Prewarm must not change the dataset version") + assertEquals(warm[0][3], retry[0][3], "Retry must cover the same eligible backend set") + + // Binary PREPARE and EXECUTE must agree on the result schema, including repeated executions. + // Force server preparation even when the JDBC URL defaults to client-side emulation. + def prepared = context.getConnection().unwrap(com.mysql.cj.jdbc.JdbcConnection) + .serverPrepareStatement(statement) + try { + assertTrue(prepared instanceof com.mysql.cj.jdbc.ServerPreparedStatement) + assertEquals(5, prepared.getMetaData().getColumnCount()) + assertEquals("DatasetVersion", prepared.getMetaData().getColumnName(3)) + for (int execution = 0; execution < 2; execution++) { + def result = exec(prepared) + assertEquals(1, result.size()) + assertEquals(indexName, result[0][1]) + assertEquals(warm[0][2].toString(), result[0][2].toString()) + assertEquals(warm[0][3].toString(), result[0][3].toString()) + } + } finally { + prepared.close() + } + + test { + sql "WARM UP INDEX missing_index ON ${tableName}" + exception "does not exist in dataset version" + } + if (isCloudMode()) { + def groups = sql "SHOW CLUSTERS" + assertFalse(groups.isEmpty()) + def explicit = sql "${statement} WITH COMPUTE GROUP `${groups[0][0]}`" + assertTrue(explicit[0][3].toInteger() > 0) + } + + // Fail one target only: successful prewarm on other BEs must not hide the failure. + def backends = sql_return_maparray "SHOW BACKENDS" + def target = backends.find { + def status = new JsonSlurper().parseText(it.Status) + it.Alive.toString().equalsIgnoreCase("true") && + !it.SystemDecommissioned.toString().equalsIgnoreCase("true") && + !status.isQueryDisabled && !status.isLoadDisabled + } + assertNotNull(target) + String failureStatement = statement + if (isCloudMode()) { + String targetGroup = new JsonSlurper().parseText(target.Tag).cloud_cluster_name Review Comment: [P2] Read the displayed compute group tag in the cloud failure test. `SHOW BACKENDS` exposes this field as `compute_group_name` (`Backend.getTagMapString()` renames it), so `cloud_cluster_name` is null here. The test sends a compute-group clause with a null name and fails before reaching the injected BE prewarm error; use `compute_group_name`. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
