github-actions[bot] commented on code in PR #68631:
URL: https://github.com/apache/doris/pull/68631#discussion_r4132497435
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducer.java:
##########
@@ -318,18 +327,63 @@ public FlightInfo getFlightInfoStatement(final
CommandStatementQuery request, fi
@Override
public FlightInfo getFlightInfoPreparedStatement(final
CommandPreparedStatementQuery command,
final CallContext context, final FlightDescriptor descriptor) {
- String[] handleParts =
command.getPreparedStatementHandle().toStringUtf8().split(":");
- String executedPeerIdentity = handleParts[0];
- String preparedStatementId = handleParts[1];
- ConnectContext connectContext =
flightSessionsManager.getConnectContext(executedPeerIdentity);
- return executeQueryStatement(executedPeerIdentity, connectContext,
- connectContext.getPreparedQuery(preparedStatementId),
descriptor);
+ ConnectContext connection =
flightSessionsManager.getConnectContext(context.peerIdentity());
+ synchronized (connection) {
+ return executeQueryStatement(context.peerIdentity(), connection,
+ preparedQuery(connection, context, command).getLeft(),
descriptor);
Review Comment:
[P2] Check the actual prepared execution schema after replanning. For
`Prepare("SELECT * FROM t")`, another session can add a column after
`preparedQuery()` reanalyzes `t` but before `executeQueryStatement()` replans
it; the session monitor does not block that DDL. GetFlightInfo then returns the
new BE schema while Prepare already advertised the old dataset schema. Compare
the returned FlightInfo schema with the stored prepared schema and expire the
handle and clean up the just-started query on mismatch.
##########
fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java:
##########
@@ -907,15 +924,34 @@ public void resetLoginTime() {
this.loginTime = System.currentTimeMillis();
}
- public void addPreparedQuery(String preparedStatementId, String
preparedQuery) {
- preparedQuerys.put(preparedStatementId, preparedQuery);
+ public synchronized void addPreparedQuery(String preparedStatementId,
String preparedQuery) {
+ addPreparedQuery(preparedStatementId, preparedQuery, null);
+ }
+
+ public synchronized void addPreparedQuery(String preparedStatementId,
String preparedQuery, Schema schema) {
+ preparedQuerys.put(preparedStatementId,
+ new PreparedQuery(preparedQuery, getDefaultCatalog(),
getDatabase(), schema));
}
- public String getPreparedQuery(String preparedStatementId) {
- return preparedQuerys.get(preparedStatementId);
+ public synchronized Schema getPreparedQuerySchema(String
preparedStatementId) {
+ PreparedQuery query = preparedQuerys.get(preparedStatementId);
+ return query == null ? null : query.schema;
+ }
+
+ public synchronized String getPreparedQuery(String preparedStatementId) {
+ PreparedQuery query = preparedQuerys.get(preparedStatementId);
+ if (query == null) {
+ return null;
+ }
+ // A handle must not execute unqualified SQL in a different namespace
than its advertised schema.
+ if (!Objects.equals(query.catalog, getDefaultCatalog()) ||
!Objects.equals(query.database, getDatabase())) {
Review Comment:
[P2] Keep explicitly namespace-changing prepared statements reusable. A
handle prepared for `USE db2; SELECT ...` in db1 runs successfully once,
leaving the session in db2. On the next execution this lookup removes the
handle because Prepare stored db1, even though the statement itself selects db2
again and has the same result schema. A standalone prepared `USE db2` also
expires after its first run. Account for the statement's own namespace
transition when validating its handle.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java:
##########
@@ -0,0 +1,364 @@
+// 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.service.arrowflight;
+
+import org.apache.doris.analysis.StatementBase;
+import org.apache.doris.catalog.AggStateType;
+import org.apache.doris.catalog.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.es.EsExternalCatalog;
+import org.apache.doris.datasource.lance.LanceExternalCatalog;
+import org.apache.doris.mysql.MysqlCommand;
+import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.CascadesContext;
+import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.glue.LogicalPlanAdapter;
+import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.parser.SqlDialectHelper;
+import org.apache.doris.nereids.rules.rewrite.CheckPrivileges;
+import org.apache.doris.nereids.trees.expressions.Slot;
+import org.apache.doris.nereids.trees.plans.Plan;
+import org.apache.doris.nereids.trees.plans.PrepareCommandPlanner;
+import org.apache.doris.nereids.trees.plans.commands.AlterTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.Command;
+import org.apache.doris.nereids.trees.plans.commands.DeleteFromCommand;
+import org.apache.doris.nereids.trees.plans.commands.DescribeCommand;
+import org.apache.doris.nereids.trees.plans.commands.ExplainCommand;
+import org.apache.doris.nereids.trees.plans.commands.KillCommand;
+import org.apache.doris.nereids.trees.plans.commands.ReplayCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowCreateTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowDataCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowPartitionsCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowProcCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowPythonPackagesCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowQueryStatsCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.TransactionCommand;
+import org.apache.doris.nereids.trees.plans.commands.UpdateCommand;
+import org.apache.doris.nereids.trees.plans.commands.info.CreateIndexOp;
+import org.apache.doris.nereids.trees.plans.commands.info.DropIndexOp;
+import
org.apache.doris.nereids.trees.plans.commands.insert.BatchInsertIntoTableCommand;
+import
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTVFCommand;
+import
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand;
+import
org.apache.doris.nereids.trees.plans.commands.insert.InsertOverwriteTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.merge.MergeIntoCommand;
+import org.apache.doris.nereids.trees.plans.commands.use.SwitchCommand;
+import org.apache.doris.nereids.trees.plans.commands.use.UseCommand;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.QueryState;
+import org.apache.doris.qe.ResultSetMetaData;
+import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.qe.ShowResultSetMetaData;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.qe.VariableMgr;
+
+import org.apache.arrow.flight.CallStatus;
+import org.apache.arrow.util.AutoCloseables;
+import org.apache.arrow.vector.types.pojo.ArrowType;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/** Resolves result metadata without scheduling fragments or evaluating query
expressions. */
+final class FlightSqlQuerySchema {
+ private FlightSqlQuerySchema() {
+ }
+
+ static Schema analyze(ConnectContext context, String query) throws
Exception {
+ synchronized (context) {
+ ConnectContext previousThreadContext = ConnectContext.get();
+ StatementContext previousStatement = context.getStatementContext();
+ SessionVariable previousSession = context.getSessionVariable();
+ QueryState previousState = context.getState();
+ StmtExecutor previousExecutor = context.getExecutor();
+ String previousCatalog = context.getDefaultCatalog();
+ String previousDatabase = context.getDatabase();
+ List<StatementBase> statements = Collections.emptyList();
+ try {
+ context.setThreadLocalInfo();
+ context.setCommand(MysqlCommand.COM_QUERY);
+ // Parsing SET_VAR hints already mutates session variables.
Isolate them even when parsing fails.
+
context.setSessionVariable(VariableMgr.cloneSessionVariable(previousSession));
+ context.setState(new QueryState());
+ context.setExecutor(null);
+ context.setStatementContext(null);
+ // Match execution's HTTP/plugin conversion before the dialect
parser sees the SQL.
+ String converted = SqlDialectHelper.convertSqlByDialect(query,
context.getSessionVariable());
+ try {
+ statements = new NereidsParser().parseSQL(converted,
context.getSessionVariable());
+ } catch (Exception convertedError) {
+ if
(!context.getSessionVariable().isRetryOriginSqlOnConvertFail() ||
converted.equals(query)) {
+ throw convertedError;
+ }
+ // Match execution's parse fallback while discarding any
failed parser context.
+ StatementContext failed = context.getStatementContext();
+ if (failed != null) {
+ failed.close();
+ context.setStatementContext(null);
+ }
+ statements = new NereidsParser().parseSQL(query,
context.getSessionVariable());
+ }
+ Map<String, String> scopedDatabases = new HashMap<>();
+ if (statements.isEmpty()) {
+ throw CallStatus.UNIMPLEMENTED.withDescription(
+ "Schema discovery requires a
statement").toRuntimeException();
+ }
+ // JDBC clients commonly prefix their query with USE. Resolve
that namespace only within this scope.
+ for (int i = 0; i < statements.size() - 1; ++i) {
+ Plan prefix = ((LogicalPlanAdapter)
statements.get(i)).getLogicalPlan();
+ if (!(prefix instanceof UseCommand) && !(prefix instanceof
SwitchCommand)) {
+ throw CallStatus.UNIMPLEMENTED.withDescription(
+ "Schema discovery only supports USE or SWITCH
before the result statement")
+ .toRuntimeException();
+ }
+ resolveNamespace(context, prefix, scopedDatabases);
+ }
+ LogicalPlanAdapter statement = (LogicalPlanAdapter)
statements.get(statements.size() - 1);
+ StatementContext statementContext =
statement.getStatementContext();
+ context.setStatementContext(statementContext);
+ statementContext.setParsedStatement(statement);
+ if (!statementContext.getPlaceholders().isEmpty()) {
+ throw CallStatus.UNIMPLEMENTED.withDescription(
+ "Flight SQL parameter binding is not
supported").toRuntimeException();
+ }
+ List<Field> fields = new ArrayList<>();
+ Plan plan = statement.getLogicalPlan();
+ if (plan instanceof Command) {
+ resolveNamespace(context, plan, scopedDatabases);
+ ResultSetMetaData metadata = commandMetadata(context,
(Command) plan);
+ if (metadata == null) {
+ throw
CallStatus.UNIMPLEMENTED.withDescription("Command result metadata is
unavailable")
+ .toRuntimeException();
+ }
+ // FE-local result sets are serialized as nullable strings
by FlightSqlChannel.
+ for (Column column : metadata.getColumns()) {
+ fields.add(Field.nullable(column.getName(), new
ArrowType.Utf8()));
+ }
+ if (fields.isEmpty()) {
+ switch (((Command) plan).stmtType()) {
+ case SET:
+ case USE:
+ case SWITCH:
+ case CREATE:
+ case ALTER:
+ case DROP:
+ case TRUNCATE:
+ // Only known no-row command categories have
the protocol OK schema.
+ fields.add(Field.nullable("StatusResult", new
ArrowType.Utf8()));
+ break;
+ default:
+ throw CallStatus.UNIMPLEMENTED.withDescription(
Review Comment:
[P2] Preserve schema discovery for HELP. `HelpCommand.getMetaData()` is
empty, so every valid `HELP ...` reaches this SHOW rejection during GetSchema
and Prepare, although direct Flight execution sends a result from
`HelpCommand.doRun()` with a TOPIC, CATEGORY, or KEYWORD header. Resolve that
header from the local help catalog without sending rows, or keep a compatible
prepared path for HELP.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java:
##########
@@ -0,0 +1,364 @@
+// 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.service.arrowflight;
+
+import org.apache.doris.analysis.StatementBase;
+import org.apache.doris.catalog.AggStateType;
+import org.apache.doris.catalog.ArrayType;
+import org.apache.doris.catalog.Column;
+import org.apache.doris.catalog.MapType;
+import org.apache.doris.catalog.PrimitiveType;
+import org.apache.doris.catalog.ScalarType;
+import org.apache.doris.catalog.StructField;
+import org.apache.doris.catalog.StructType;
+import org.apache.doris.catalog.Type;
+import org.apache.doris.datasource.CatalogIf;
+import org.apache.doris.datasource.es.EsExternalCatalog;
+import org.apache.doris.datasource.lance.LanceExternalCatalog;
+import org.apache.doris.mysql.MysqlCommand;
+import org.apache.doris.mysql.privilege.PrivPredicate;
+import org.apache.doris.nereids.CascadesContext;
+import org.apache.doris.nereids.StatementContext;
+import org.apache.doris.nereids.glue.LogicalPlanAdapter;
+import org.apache.doris.nereids.parser.NereidsParser;
+import org.apache.doris.nereids.parser.SqlDialectHelper;
+import org.apache.doris.nereids.rules.rewrite.CheckPrivileges;
+import org.apache.doris.nereids.trees.expressions.Slot;
+import org.apache.doris.nereids.trees.plans.Plan;
+import org.apache.doris.nereids.trees.plans.PrepareCommandPlanner;
+import org.apache.doris.nereids.trees.plans.commands.AlterTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.Command;
+import org.apache.doris.nereids.trees.plans.commands.DeleteFromCommand;
+import org.apache.doris.nereids.trees.plans.commands.DescribeCommand;
+import org.apache.doris.nereids.trees.plans.commands.ExplainCommand;
+import org.apache.doris.nereids.trees.plans.commands.KillCommand;
+import org.apache.doris.nereids.trees.plans.commands.ReplayCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowCreateTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowDataCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowPartitionsCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowProcCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowPythonPackagesCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowQueryStatsCommand;
+import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.TransactionCommand;
+import org.apache.doris.nereids.trees.plans.commands.UpdateCommand;
+import org.apache.doris.nereids.trees.plans.commands.info.CreateIndexOp;
+import org.apache.doris.nereids.trees.plans.commands.info.DropIndexOp;
+import
org.apache.doris.nereids.trees.plans.commands.insert.BatchInsertIntoTableCommand;
+import
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTVFCommand;
+import
org.apache.doris.nereids.trees.plans.commands.insert.InsertIntoTableCommand;
+import
org.apache.doris.nereids.trees.plans.commands.insert.InsertOverwriteTableCommand;
+import org.apache.doris.nereids.trees.plans.commands.merge.MergeIntoCommand;
+import org.apache.doris.nereids.trees.plans.commands.use.SwitchCommand;
+import org.apache.doris.nereids.trees.plans.commands.use.UseCommand;
+import org.apache.doris.qe.ConnectContext;
+import org.apache.doris.qe.QueryState;
+import org.apache.doris.qe.ResultSetMetaData;
+import org.apache.doris.qe.SessionVariable;
+import org.apache.doris.qe.ShowResultSetMetaData;
+import org.apache.doris.qe.StmtExecutor;
+import org.apache.doris.qe.VariableMgr;
+
+import org.apache.arrow.flight.CallStatus;
+import org.apache.arrow.util.AutoCloseables;
+import org.apache.arrow.vector.types.pojo.ArrowType;
+import org.apache.arrow.vector.types.pojo.Field;
+import org.apache.arrow.vector.types.pojo.FieldType;
+import org.apache.arrow.vector.types.pojo.Schema;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+/** Resolves result metadata without scheduling fragments or evaluating query
expressions. */
+final class FlightSqlQuerySchema {
+ private FlightSqlQuerySchema() {
+ }
+
+ static Schema analyze(ConnectContext context, String query) throws
Exception {
+ synchronized (context) {
+ ConnectContext previousThreadContext = ConnectContext.get();
+ StatementContext previousStatement = context.getStatementContext();
+ SessionVariable previousSession = context.getSessionVariable();
+ QueryState previousState = context.getState();
+ StmtExecutor previousExecutor = context.getExecutor();
+ String previousCatalog = context.getDefaultCatalog();
+ String previousDatabase = context.getDatabase();
+ List<StatementBase> statements = Collections.emptyList();
+ try {
+ context.setThreadLocalInfo();
+ context.setCommand(MysqlCommand.COM_QUERY);
+ // Parsing SET_VAR hints already mutates session variables.
Isolate them even when parsing fails.
+
context.setSessionVariable(VariableMgr.cloneSessionVariable(previousSession));
+ context.setState(new QueryState());
+ context.setExecutor(null);
+ context.setStatementContext(null);
+ // Match execution's HTTP/plugin conversion before the dialect
parser sees the SQL.
+ String converted = SqlDialectHelper.convertSqlByDialect(query,
context.getSessionVariable());
+ try {
+ statements = new NereidsParser().parseSQL(converted,
context.getSessionVariable());
+ } catch (Exception convertedError) {
+ if
(!context.getSessionVariable().isRetryOriginSqlOnConvertFail() ||
converted.equals(query)) {
+ throw convertedError;
+ }
+ // Match execution's parse fallback while discarding any
failed parser context.
+ StatementContext failed = context.getStatementContext();
+ if (failed != null) {
+ failed.close();
+ context.setStatementContext(null);
+ }
+ statements = new NereidsParser().parseSQL(query,
context.getSessionVariable());
+ }
+ Map<String, String> scopedDatabases = new HashMap<>();
+ if (statements.isEmpty()) {
+ throw CallStatus.UNIMPLEMENTED.withDescription(
+ "Schema discovery requires a
statement").toRuntimeException();
+ }
+ // JDBC clients commonly prefix their query with USE. Resolve
that namespace only within this scope.
+ for (int i = 0; i < statements.size() - 1; ++i) {
+ Plan prefix = ((LogicalPlanAdapter)
statements.get(i)).getLogicalPlan();
+ if (!(prefix instanceof UseCommand) && !(prefix instanceof
SwitchCommand)) {
+ throw CallStatus.UNIMPLEMENTED.withDescription(
+ "Schema discovery only supports USE or SWITCH
before the result statement")
+ .toRuntimeException();
+ }
+ resolveNamespace(context, prefix, scopedDatabases);
+ }
+ LogicalPlanAdapter statement = (LogicalPlanAdapter)
statements.get(statements.size() - 1);
+ StatementContext statementContext =
statement.getStatementContext();
+ context.setStatementContext(statementContext);
+ statementContext.setParsedStatement(statement);
+ if (!statementContext.getPlaceholders().isEmpty()) {
+ throw CallStatus.UNIMPLEMENTED.withDescription(
+ "Flight SQL parameter binding is not
supported").toRuntimeException();
+ }
+ List<Field> fields = new ArrayList<>();
+ Plan plan = statement.getLogicalPlan();
+ if (plan instanceof Command) {
+ resolveNamespace(context, plan, scopedDatabases);
+ ResultSetMetaData metadata = commandMetadata(context,
(Command) plan);
+ if (metadata == null) {
+ throw
CallStatus.UNIMPLEMENTED.withDescription("Command result metadata is
unavailable")
+ .toRuntimeException();
+ }
+ // FE-local result sets are serialized as nullable strings
by FlightSqlChannel.
+ for (Column column : metadata.getColumns()) {
+ fields.add(Field.nullable(column.getName(), new
ArrowType.Utf8()));
+ }
+ if (fields.isEmpty()) {
+ switch (((Command) plan).stmtType()) {
+ case SET:
+ case USE:
+ case SWITCH:
+ case CREATE:
+ case ALTER:
+ case DROP:
+ case TRUNCATE:
+ // Only known no-row command categories have
the protocol OK schema.
+ fields.add(Field.nullable("StatusResult", new
ArrowType.Utf8()));
+ break;
+ default:
+ throw CallStatus.UNIMPLEMENTED.withDescription(
+ "Result metadata is unavailable
without executing this command")
+ .toRuntimeException();
+ }
+ }
+ } else {
+ PrepareCommandPlanner planner = new
PrepareCommandPlanner(statementContext);
+ planner.plan(statement,
context.getSessionVariable().toThrift());
+ CascadesContext cascades = planner.getCascadesContext();
+ Plan analyzed = cascades.getRewritePlan();
+ // PrepareCommandPlanner stops before the rewrite phase
that normally checks privileges.
+ new CheckPrivileges().rewriteRoot(analyzed,
cascades.getCurrentJobContext());
+ for (Slot slot : analyzed.getOutput()) {
+ fields.add(field(slot.getName(),
slot.getDataType().toCatalogDataType(), slot.nullable(),
+ true,
context.getSessionVariable().getTimeZone()));
+ }
+ }
+ return new Schema(fields);
+ } finally {
+ try {
+ List<AutoCloseable> resources = new ArrayList<>();
+ for (StatementBase statement : statements) {
+ if (statement instanceof LogicalPlanAdapter) {
+ resources.add(((LogicalPlanAdapter)
statement).getStatementContext());
+ }
+ }
+ // A parser failure can leave a context that was never
added to the returned list.
+ StatementContext current = context.getStatementContext();
+ if (current != null && current != previousStatement &&
!resources.contains(current)) {
+ resources.add(current);
+ }
+ AutoCloseables.close(resources);
+ } finally {
+ try {
+ context.setStatementContext(previousStatement);
+ context.setSessionVariable(previousSession);
+ context.setState(previousState);
+ context.setExecutor(previousExecutor);
+ if
(!previousCatalog.equals(context.getDefaultCatalog())
+ ||
!previousDatabase.equals(context.getDatabase())) {
+ context.changeDefaultCatalog(previousCatalog);
+ context.setDatabase(previousDatabase);
+ }
+ } finally {
+ context.setCommand(MysqlCommand.COM_SLEEP);
+ if (previousThreadContext == null) {
+ ConnectContext.remove();
+ } else {
+ previousThreadContext.setThreadLocalInfo();
+ }
+ }
+ }
+ }
+ }
+ }
+
+ private static ResultSetMetaData commandMetadata(ConnectContext context,
Command command) throws Exception {
+ // These getters depend on execution-time state or remote responses.
Do not advertise a
+ // guessed schema, or run the command merely to discover it.
+ if (command instanceof ShowPythonPackagesCommand || command instanceof
DescribeCommand
+ || command instanceof ShowDataCommand || command instanceof
ShowPartitionsCommand
+ || command instanceof ShowQueryStatsCommand) {
+ throw CallStatus.UNIMPLEMENTED.withDescription("Command schema
requires execution-time metadata")
+ .toRuntimeException();
+ }
+ if (command instanceof ShowTableCommand) {
+ ((ShowTableCommand) command).validate(context);
+ } else if (command instanceof ShowCreateTableCommand) {
+ return ((ShowCreateTableCommand) command).getMetaData(context);
+ } else if (command instanceof ShowProcCommand) {
+ return ((ShowProcCommand) command).getMetaData(context);
+ }
+ if (command instanceof ExplainCommand) {
+ // PLAN PROCESS has no Flight serialization path in StmtExecutor.
+ if (((ExplainCommand) command).showPlanProcess()) {
+ throw CallStatus.UNIMPLEMENTED.withDescription("EXPLAIN PLAN
PROCESS is not supported over Flight SQL")
+ .toRuntimeException();
+ }
+ return stringMetadata("Explain String(Nereids Planner)");
+ } else if (command instanceof ReplayCommand) {
+ return stringMetadata("Plan Replayer dump url");
+ } else if (command instanceof AlterTableCommand) {
+ AlterTableCommand alter = (AlterTableCommand) command;
+ String catalog = alter.getTbl().getCtl();
+ // Lance index admission returns a JobId header even for an IF
no-op. Do not run
+ // validation/admission here: those paths can resolve remote
tables or allocate IDs.
+ if (context.getCatalog(catalog == null ?
context.getDefaultCatalog() : catalog)
+ instanceof LanceExternalCatalog &&
alter.getNereidsOps().stream().anyMatch(op ->
+ (op instanceof CreateIndexOp && !((CreateIndexOp)
op).isAlter())
+ || (op instanceof DropIndexOp &&
!((DropIndexOp) op).isAlter()))) {
+ return stringMetadata("JobId");
+ }
+ }
+ ResultSetMetaData metadata = command.getResultSetMetaData();
Review Comment:
[P2] Resolve the filtered SHOW SNAPSHOT header before advertising it. For
`SHOW SNAPSHOT ON repo WHERE SNAPSHOT = 's' AND TIMESTAMP = 't'`, this getter
runs before `ShowSnapshotCommand.validate()` fills `snapshotName` and
`timestamp`, so Prepare/GetSchema advertises the three-column SNAPSHOT_ALL
header. Execution then uses the five-column SNAPSHOT_DETAIL header and rows.
Determine the header from the predicates without fetching snapshots, or decline
schema discovery for this variant.
--
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]