github-actions[bot] commented on code in PR #68631: URL: https://github.com/apache/doris/pull/68631#discussion_r4131745638
########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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); + statements = new NereidsParser().parseSQL(query, context.getSessionVariable()); Review Comment: [P2] Apply dialect conversion before schema parsing. `ConnectProcessor.executeQuery` first runs `SqlDialectHelper.convertSqlByDialect`, but this path parses the original text. For a Flight session using a non-Doris dialect and the built-in HTTP converter, `parseSqlWithDialect` returns null, so a query whose syntax needs conversion executes through GetFlightInfo but fails GetSchema and prepared-statement creation. Use the same conversion path before analysis. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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); + statements = new NereidsParser().parseSQL(query, context.getSessionVariable()); + 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); + } + 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); + if (plan instanceof ShowTableCommand) { + // SHOW TABLES labels include the database normally resolved when the command runs. + ((ShowTableCommand) plan).validate(context); + } + ResultSetMetaData metadata = ((Command) plan).getResultSetMetaData(); + 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()) { Review Comment: [P2] Do not use empty command metadata as an OK-result schema. The base `Command.getResultSetMetaData()` returns an empty list even for result-producing commands. For example, `SHOW FRONTEND CONFIG` inherits that default and `StmtType.OTHER`, so GetSchema and Prepare advertise one `StatusResult` field, while execution sends six configuration columns; COPY INTO and WARM UP CLUSTER have the same problem. Return their result metadata or leave schema discovery unsupported for them. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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)); Review Comment: [P1] Isolate schema analysis from concurrent Flight execution. `synchronized(context)` only serializes schema calls; `getFlightInfoStatement` executes on the same context without that lock. If a GetSchema/Prepare call installs the cloned `SessionVariable` here while another RPC executes `SET query_timeout=17`, that SET can complete successfully on the clone, then this method restores the old variable and silently loses the change. Serialize all per-session execution with schema analysis or analyze in an independent context. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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); + statements = new NereidsParser().parseSQL(query, context.getSessionVariable()); + 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); + } + 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); + if (plan instanceof ShowTableCommand) { + // SHOW TABLES labels include the database normally resolved when the command runs. + ((ShowTableCommand) plan).validate(context); + } + ResultSetMetaData metadata = ((Command) plan).getResultSetMetaData(); + 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 SHOW: + case EXPLAIN: + case CALL: + case EXECUTE: + case PREPARE: + throw CallStatus.UNIMPLEMENTED.withDescription( + "Result metadata is unavailable without executing this command") + .toRuntimeException(); + default: + // Commands without rows use this actual protocol result in executeQueryStatement. + // Keeping it nonempty also prevents JDBC from selecting the unsupported update RPC. + fields.add(Field.nullable("StatusResult", new ArrowType.Utf8())); + } + } + } 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 void resolveNamespace(ConnectContext context, Plan plan) throws Exception { + if (plan instanceof UseCommand) { + UseCommand use = (UseCommand) plan; + String catalog = use.getCatalogName() == null ? context.getDefaultCatalog() : use.getCatalogName(); + CatalogIf catalogObject = context.getCatalog(catalog); + if (catalogObject == null || !context.getEnv().getAccessManager() + .checkDbPriv(context, catalog, use.getDatabaseName(), PrivPredicate.SHOW)) { + throw CallStatus.UNAUTHORIZED.withDescription("Database access denied").toRuntimeException(); + } + catalogObject.getDbOrAnalysisException(use.getDatabaseName()); + context.changeDefaultCatalog(catalog); + context.setDatabase(use.getDatabaseName()); + } else if (plan instanceof SwitchCommand) { + String catalog = ((SwitchCommand) plan).getCatalogName(); + if (context.getCatalog(catalog) == null || !context.getEnv().getAccessManager() + .checkCtlPriv(context, catalog, PrivPredicate.SHOW)) { + throw CallStatus.UNAUTHORIZED.withDescription("Catalog access denied").toRuntimeException(); + } + context.changeDefaultCatalog(catalog); Review Comment: [P2] Restore the target database when analyzing SWITCH. `changeDefaultCatalog` clears the current database, while executed `SWITCH` calls `Env.changeCatalog`, which restores that catalog's remembered database (or the ES default). After a session has used catalog B/database b and returned to A, `SWITCH B; SELECT * FROM t` executes against B.b but GetSchema/Prepare analyze with an empty database and fail. Reproduce the execution namespace transition within this scoped analysis. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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); + statements = new NereidsParser().parseSQL(query, context.getSessionVariable()); + 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); + } + 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); + if (plan instanceof ShowTableCommand) { + // SHOW TABLES labels include the database normally resolved when the command runs. + ((ShowTableCommand) plan).validate(context); + } + ResultSetMetaData metadata = ((Command) plan).getResultSetMetaData(); Review Comment: [P2] Resolve command-dependent metadata before returning a schema. This getter runs before execution determines the result shape: `SHOW PYTHON PACKAGES` advertises two fields but can return four, OLAP `DESCRIBE t ALL` advertises six MySQL fields but returns 12, and `SHOW DATA FROM table` advertises database-wide fields until `validate()` sets its database. GetSchema and prepared metadata can therefore disagree with execution. Determine the actual shape first or decline schema discovery for these commands. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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); + statements = new NereidsParser().parseSQL(query, context.getSessionVariable()); + 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); + } + 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); + if (plan instanceof ShowTableCommand) { Review Comment: [P2] Initialize metadata-dependent SHOW commands before reading their schema. This path validates only `ShowTableCommand`; `ShowPartitionsCommand.getMetaData()` dereferences `catalog` and `ShowQueryStatsCommand.getMetaData()` switches on `type`, both of which are set only by their own `validate(ctx)` methods. Their valid statements execute after validation but GetSchema/Prepare throw before returning metadata. Perform the required command setup or report metadata as unavailable. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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); + statements = new NereidsParser().parseSQL(query, context.getSessionVariable()); + 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); + } + 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); + if (plan instanceof ShowTableCommand) { + // SHOW TABLES labels include the database normally resolved when the command runs. + ((ShowTableCommand) plan).validate(context); + } + ResultSetMetaData metadata = ((Command) plan).getResultSetMetaData(); + 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 SHOW: Review Comment: [P2] Preserve preparation for executable SHOW commands with dynamic metadata. `SHOW CREATE TABLE t` returns an empty `getMetaData()` here, so this SHOW branch throws UNIMPLEMENTED, even though `doRun()` returns a nonempty table/view schema. Because `createPreparedStatement` now calls this analysis before registering a handle, previously executable prepared SHOW CREATE TABLE and SHOW PROC requests fail at Prepare. Resolve their metadata with the required checks or preserve a compatible preparation path. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java: ########## @@ -0,0 +1,250 @@ +// 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.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.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.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.Command; +import org.apache.doris.nereids.trees.plans.commands.ShowTableCommand; +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.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.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); + statements = new NereidsParser().parseSQL(query, context.getSessionVariable()); + 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); + } + 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); + if (plan instanceof ShowTableCommand) { + // SHOW TABLES labels include the database normally resolved when the command runs. + ((ShowTableCommand) plan).validate(context); + } + ResultSetMetaData metadata = ((Command) plan).getResultSetMetaData(); + 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 SHOW: + case EXPLAIN: + case CALL: + case EXECUTE: + case PREPARE: + throw CallStatus.UNIMPLEMENTED.withDescription( + "Result metadata is unavailable without executing this command") + .toRuntimeException(); + default: + // Commands without rows use this actual protocol result in executeQueryStatement. + // Keeping it nonempty also prevents JDBC from selecting the unsupported update RPC. + fields.add(Field.nullable("StatusResult", new ArrowType.Utf8())); + } + } + } 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 void resolveNamespace(ConnectContext context, Plan plan) throws Exception { + if (plan instanceof UseCommand) { + UseCommand use = (UseCommand) plan; + String catalog = use.getCatalogName() == null ? context.getDefaultCatalog() : use.getCatalogName(); + CatalogIf catalogObject = context.getCatalog(catalog); + if (catalogObject == null || !context.getEnv().getAccessManager() + .checkDbPriv(context, catalog, use.getDatabaseName(), PrivPredicate.SHOW)) { + throw CallStatus.UNAUTHORIZED.withDescription("Database access denied").toRuntimeException(); + } + catalogObject.getDbOrAnalysisException(use.getDatabaseName()); + context.changeDefaultCatalog(catalog); + context.setDatabase(use.getDatabaseName()); + } else if (plan instanceof SwitchCommand) { + String catalog = ((SwitchCommand) plan).getCatalogName(); + if (context.getCatalog(catalog) == null || !context.getEnv().getAccessManager() + .checkCtlPriv(context, catalog, PrivPredicate.SHOW)) { + throw CallStatus.UNAUTHORIZED.withDescription("Catalog access denied").toRuntimeException(); + } + context.changeDefaultCatalog(catalog); + } + } + + private static Field field(String name, Type type, boolean nullable, boolean topLevel, String timezone) { + PrimitiveType primitive = type.getPrimitiveType(); + int precision = type instanceof ScalarType ? ((ScalarType) type).getScalarPrecision() : 0; + int scale = type instanceof ScalarType ? ((ScalarType) type).getScalarScale() : 0; + ArrowType arrowType = FlightSqlSchemaHelper.getArrowType(primitive, precision, scale); + if (primitive == PrimitiveType.TIMESTAMPTZ) { + arrowType = new ArrowType.Timestamp(((ArrowType.Timestamp) arrowType).getUnit(), + "Z".equals(timezone) ? "UTC" : timezone); + } + if (arrowType instanceof ArrowType.Null && primitive != PrimitiveType.NULL_TYPE) { Review Comment: [P2] Map aggregate states using their serialized Arrow type. `SELECT group_concat_state('x') AS s` has a non-null `AGG_STATE` output, so the helper returns Arrow Null and this branch rejects GetSchema/Prepare. The BE first unwraps `DataTypeAggState` to the aggregate's serialized type; group_concat uses string serialization and executes with an Arrow UTF8 schema. Preserve that supported execution path in schema discovery. ########## fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/DorisFlightSqlProducer.java: ########## @@ -355,51 +383,30 @@ private ActionCreatePreparedStatementResult buildCreatePreparedStatementResult(B @Override public void createPreparedStatement(final ActionCreatePreparedStatementRequest request, final CallContext context, final StreamListener<Result> listener) { - // TODO can only execute complete SQL, not support SQL parameters. - // For Python: the Python code will try to create a prepared statement (this is to fit DBAPI, IIRC) and - // if the server raises any error except for NotImplemented it will fail. (If it gets NotImplemented, - // it will ignore and execute without a prepared statement.) see: https://github.com/apache/arrow/issues/38786 executorService.submit(() -> { - ConnectContext connectContext = flightSessionsManager.getConnectContext(context.peerIdentity()); + ConnectContext connectContext = null; + String preparedStatementId = null; try { - connectContext.setCommand(MysqlCommand.COM_QUERY); - final String query = request.getQuery(); - String preparedStatementId = UUID.randomUUID().toString(); - final ByteString handle = ByteString.copyFromUtf8(context.peerIdentity() + ":" + preparedStatementId); + connectContext = flightSessionsManager.getConnectContext(context.peerIdentity()); + String query = request.getQuery(); + // ADBC ExecuteSchema reads this dataset schema directly without calling GetSchema. + // Analyze before registering a handle so failed preparation does not retain a query. + Schema schema = analyzeQuerySchema(connectContext, query); Review Comment: [P2] Keep a prepared handle's dataset schema aligned with execution. This embeds an analyzed schema in the Prepare result but registers only the unqualified SQL text. If a session prepares `SELECT * FROM t` in database A and then executes `USE B` where `t` has different columns, the prepared result still advertises A's fields while `getSchemaPreparedStatement` and execution resolve B's fields. Bind the handle to its preparation namespace or update/invalidate its advertised schema when the namespace changes. -- 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]
