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]

Reply via email to