github-actions[bot] commented on code in PR #68775:
URL: https://github.com/apache/doris/pull/68775#discussion_r4216115261


##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/WarmUpIndexCommand.java:
##########
@@ -0,0 +1,54 @@
+// 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.doris.nereids.trees.plans.commands;
+
+import org.apache.doris.datasource.lance.LanceIndexPrewarm;
+import org.apache.doris.info.TableNameInfo;
+import org.apache.doris.nereids.trees.plans.PlanType;
+import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.StmtExecutor;
+
+/** Synchronously fills the query sessions on the selected backends. */
+public class WarmUpIndexCommand extends Command {
+    private final TableNameInfo table;
+    private final String indexName;
+    private final String computeGroup;
+    private volatile boolean cancelled;
+
+    public WarmUpIndexCommand(TableNameInfo table, String indexName, String 
computeGroup) {
+        super(PlanType.WARM_UP_INDEX_COMMAND);
+        this.table = table;
+        this.indexName = indexName;
+        this.computeGroup = computeGroup;
+    }
+
+    @Override
+    public void run(ConnectContext ctx, StmtExecutor executor) throws 
Exception {

Review Comment:
   [P2] Expose the prewarm result columns to prepared clients. 
`COM_STMT_PREPARE` wraps this command in `PrepareCommand`, which reads 
`Command.getResultSetMetaData()`; because this class inherits the empty 
default, the prepare response advertises zero columns. `COM_STMT_EXECUTE` then 
runs this method and sends the five columns in `LanceIndexPrewarm.RESULT_META`, 
so binary prepared clients receive a result schema different from the one they 
prepared. Override the metadata hook with the same five columns and cover the 
binary prepare/execute path.



##########
be/src/service/internal_service.cpp:
##########
@@ -811,6 +813,60 @@ void 
PInternalService::outfile_write_success(google::protobuf::RpcController* co
     }
 }
 
+void PInternalService::prewarm_lance_index(google::protobuf::RpcController* 
controller,
+                                           const PLanceIndexPrewarmRequest* 
request,
+                                           PLanceIndexPrewarmResponse* 
response,
+                                           google::protobuf::Closure* done) {
+    const auto received = std::chrono::steady_clock::now();
+    bool offered = _heavy_work_pool.try_offer([request, response, done, 
received]() {
+        brpc::ClosureGuard closure_guard(done);
+        Status status;
+        try {
+            auto run = [&]() -> Status {
+                DBUG_EXECUTE_IF("PInternalService.prewarm_lance_index.fail", {
+                    return Status::InternalError("Injected Lance index prewarm 
failure");
+                });
+                const auto queued_ms = 
std::chrono::duration_cast<std::chrono::milliseconds>(
+                                               
std::chrono::steady_clock::now() - received)
+                                               .count();
+                if (!request->has_timeout_ms() || request->timeout_ms() <= 
queued_ms) {
+                    return Status::TimedOut("Lance index prewarm expired 
before execution");
+                }
+                // C strings cannot represent embedded NUL bytes; reject 
rather than silently
+                // opening a different path/index or truncating a vended 
credential.
+                if (request->dataset_uri().find('\0') != std::string::npos ||
+                    request->index_name().find('\0') != std::string::npos) {
+                    return Status::InvalidArgument("Invalid Lance prewarm URI 
or index name");
+                }
+                std::vector<const char*> options;
+                options.reserve(request->storage_options_size() * 2 + 1);
+                for (const auto& [key, value] : request->storage_options()) {
+                    if (key.find('\0') != std::string::npos ||
+                        value.find('\0') != std::string::npos) {
+                        return Status::InvalidArgument("Invalid Lance prewarm 
storage option");
+                    }
+                    options.push_back(key.c_str());
+                    options.push_back(value.c_str());
+                }
+                options.push_back(nullptr);
+                
RETURN_IF_ERROR(format::lance::LanceSessionManager::instance().prewarm_index(

Review Comment:
   [P2] Bound prewarm work after it leaves the RPC queue. `timeout_ms` is 
checked only against queue time before this synchronous SDK call. If 
object-store open or index prewarm runs beyond the FE deadline (or the user 
cancels), the FE returns but this heavy-pool worker and `LanceSessionManager`'s 
process-wide prewarm mutex remain occupied; every retry on this BE is rejected 
until the old call finishes. Propagate a deadline/cancellation to the SDK work 
or isolate and bound abandoned work so a timed-out request cannot monopolize 
this path.



##########
fe/fe-core/src/main/java/org/apache/doris/datasource/lance/LanceIndexPrewarm.java:
##########
@@ -0,0 +1,226 @@
+// 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.doris.datasource.lance;
+
+import org.apache.doris.analysis.ResourceTypeEnum;
+import org.apache.doris.analysis.TableName;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.Env;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.TableIf;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.UserException;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.lance.metadata.LanceTableMetadata;
+import org.apache.doris.info.TableNameInfo;
+import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.proto.InternalService.PLanceIndexPrewarmRequest;
+import org.apache.doris.proto.InternalService.PLanceIndexPrewarmResponse;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.ShowResultSet;
+import org.apache.doris.qe.ShowResultSetMetaData;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.resource.computegroup.ComputeGroup;
+import org.apache.doris.rpc.BackendServiceProxy;
+import org.apache.doris.system.Backend;
+import org.apache.doris.system.BeSelectionPolicy;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TStatusCode;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.function.BooleanSupplier;
+
+/** Pins table access and coordinates bounded synchronous prewarm RPCs. */
+public final class LanceIndexPrewarm {
+    static final int MAX_IN_FLIGHT = 8;
+    private static final ShowResultSetMetaData RESULT_META = 
ShowResultSetMetaData.builder()
+            .addColumn(new Column("Table", ScalarType.createStringType()))
+            .addColumn(new Column("Index", ScalarType.createStringType()))
+            .addColumn(new Column("DatasetVersion", ScalarType.BIGINT))
+            .addColumn(new Column("BackendCount", ScalarType.INT))
+            .addColumn(new Column("ElapsedMs", ScalarType.BIGINT)).build();
+
+    private LanceIndexPrewarm() {
+    }
+
+    public static void run(ConnectContext ctx, StmtExecutor executor, 
TableNameInfo tableName,
+            String indexName, String computeGroup, BooleanSupplier cancelled) 
throws Exception {
+        long started = System.nanoTime();
+        long elapsedMs = ctx.getStartTime() > 0 ? Math.max(0, 
System.currentTimeMillis() - ctx.getStartTime()) : 0;
+        long remainingMs = TimeUnit.SECONDS.toMillis(ctx.getQueryTimeoutS()) - 
elapsedMs;
+        long deadline = started + TimeUnit.MILLISECONDS.toNanos(Math.max(0, 
remainingMs));
+        tableName.analyze(ctx);
+        checkPrivileges(ctx, tableName);
+        List<Backend> targets = selectBackends(resolveComputeGroup(ctx, 
computeGroup));
+        CatalogIf<?> catalog = 
Env.getCurrentEnv().getCatalogMgr().getCatalog(tableName.getCtl());
+        if (!(catalog instanceof LanceExternalCatalog)) {
+            throw new AnalysisException("WARM UP INDEX requires a Lance 
catalog table");
+        }
+        TableIf table = catalog.getDbOrAnalysisException(tableName.getDb())
+                .getTableOrAnalysisException(tableName.getTbl());
+        if (!(table instanceof LanceExternalTable)) {
+            throw new AnalysisException("WARM UP INDEX requires a Lance 
catalog table");
+        }
+        checkActive(deadline, cancelled);
+        // This is the same snapshot/access path used by queries, including 
REST-vended credentials.
+        // Never independently resolve latest, index segments, or credentials 
on each backend.
+        LanceTableMetadata metadata = ((LanceExternalTable) 
table).loadMetadata();

Review Comment:
   [P2] Include metadata loading in the command's timeout and cancellation 
path. `loadMetadata()` synchronously resolves Lance namespace access and 
opens/reads the remote dataset, but the next `checkActive` is only after that 
call returns. With a slow object store, `KILL QUERY` or `query_timeout` leaves 
the FE command thread blocked well past its deadline before any BE RPC is sent. 
Run this phase with a bounded, cancellable deadline and stop before dispatch if 
it expires.



##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/WarmUpIndexCommand.java:
##########
@@ -0,0 +1,54 @@
+// 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.doris.nereids.trees.plans.commands;
+
+import org.apache.doris.datasource.lance.LanceIndexPrewarm;
+import org.apache.doris.info.TableNameInfo;
+import org.apache.doris.nereids.trees.plans.PlanType;
+import org.apache.doris.nereids.trees.plans.visitor.PlanVisitor;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.StmtExecutor;
+
+/** Synchronously fills the query sessions on the selected backends. */
+public class WarmUpIndexCommand extends Command {
+    private final TableNameInfo table;
+    private final String indexName;
+    private final String computeGroup;
+    private volatile boolean cancelled;
+
+    public WarmUpIndexCommand(TableNameInfo table, String indexName, String 
computeGroup) {
+        super(PlanType.WARM_UP_INDEX_COMMAND);
+        this.table = table;
+        this.indexName = indexName;
+        this.computeGroup = computeGroup;
+    }
+
+    @Override
+    public void run(ConnectContext ctx, StmtExecutor executor) throws 
Exception {
+        LanceIndexPrewarm.run(ctx, executor, table, indexName, computeGroup, 
() -> cancelled || ctx.isKilled());
+    }
+
+    public void cancel() {

Review Comment:
   [P2] Keep cancellation state scoped to one execution. `ExecuteCommand.run` 
reuses the `WarmUpIndexCommand` stored in `PrepareCommand` for each 
`COM_STMT_EXECUTE`; after `KILL QUERY` calls this setter, `cancelled` stays 
true on that retained object. A subsequent execute on the same live connection 
reaches `LanceIndexPrewarm.checkActive` and immediately fails as cancelled. 
Reset or replace per-execution state without losing cancellation of the active 
run.



-- 
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