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]
