Gabriel39 commented on code in PR #68631: URL: https://github.com/apache/doris/pull/68631#discussion_r4131945305
########## 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: Addressed in a02347b241d. Added context-aware metadata resolution for SHOW CREATE TABLE and SHOW PROC. SHOW CREATE TABLE validates access and resolves the table/view kind without rendering DDL; SHOW PROC validates ADMIN/NODE access and resolves the proc node's header. Both commands are covered by prepared-schema tests, and SHOW PROC also has Flight RPC coverage. ########## 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: Addressed in a02347b241d. Prepared entries now remember their preparation catalog/database. A lookup from a different namespace removes the handle and returns NOT_FOUND with an instruction to prepare again; validation and execution are atomic under the session lock. Both prepared GetSchema and GetFlightInfo reject the expired handle. Added FE and Flight RPC regressions that change the database after Prepare. -- 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]
