Gabriel39 commented on code in PR #68631:
URL: https://github.com/apache/doris/pull/68631#discussion_r4131943416


##########
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:
   Addressed in a02347b241d. Normal and prepared execution now share the 
session context monitor with schema analysis. Preparation also keeps analysis 
and handle registration atomic. A two-worker regression pauses analysis after 
installing its clone, attempts SET concurrently, and verifies that execution 
waits and the successful SET updates the original session. Before the fix the 
final timeout reverted to 900 instead of 17; the test now passes.



##########
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:
   Addressed in a02347b241d. Mapped group_concat aggregate states to their BE 
string serialization (Arrow UTF8). Added Prepare/GetSchema parity in FE tests 
and the persistent Flight RPC regression. This intentionally does not guess the 
wire types of other aggregate states: unmodeled serializations, including 
fixed-size states, remain explicitly unsupported.



##########
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:
   Addressed in a02347b241d. SHOW PARTITIONS and SHOW QUERY STATS now return 
UNIMPLEMENTED before accessing their uninitialized metadata getters. Added 
coverage for both GetSchema and Prepare. This follows the explicit-unavailable 
option: schema discovery does not execute the commands or return a misleading 
header.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to