github-actions[bot] commented on code in PR #68631:
URL: https://github.com/apache/doris/pull/68631#discussion_r4132177975
##########
fe/fe-core/src/main/java/org/apache/doris/qe/ConnectContext.java:
##########
@@ -907,15 +921,24 @@ 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) {
+ preparedQuerys.put(preparedStatementId, new
PreparedQuery(preparedQuery, getDefaultCatalog(), getDatabase()));
}
- public String getPreparedQuery(String preparedStatementId) {
- return preparedQuerys.get(preparedStatementId);
+ 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] Revalidate the prepared schema after result-affecting session changes.
This lookup checks only catalog and database, so `Prepare("SELECT 1 || 2 AS
x")` under the default `sql_mode` advertises Bool, then `SET
sql_mode='PIPES_AS_CONCAT'` leaves the handle valid while prepared execution
reparses `||` as string concatenation and returns UTF8. This differs from the
resolved namespace thread because the namespace never changes. Capture or
recheck schema-affecting settings before GetSchema and GetFlightInfo, or
invalidate the handle when they change.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java:
##########
@@ -0,0 +1,301 @@
+// 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.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.Command;
+import org.apache.doris.nereids.trees.plans.commands.DescribeCommand;
+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.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.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());
+ statements = new NereidsParser().parseSQL(converted,
context.getSessionVariable());
Review Comment:
[P2] Match execution's parse fallback after dialect conversion. When
`retry_origin_sql_on_convert_fail` is enabled,
`ConnectProcessor.parseWithFallback` retries the original SQL if a dialect
plugin returns nonempty but unparseable converted text. This new path parses
only the converted text, so a valid original `SELECT` can execute through
GetFlightInfo while GetSchema and Prepare fail. Apply the same scoped fallback
here; this is separate from the earlier missing-conversion thread.
##########
fe/fe-core/src/main/java/org/apache/doris/nereids/trees/plans/commands/ShowProcCommand.java:
##########
@@ -74,6 +74,15 @@ public ShowResultSetMetaData getMetaData() {
return ShowResultSetMetaData.builder().build();
}
+ /** Resolve a proc node's header with the same privilege checks as SHOW
PROC. */
+ public ShowResultSetMetaData getMetaData(ConnectContext ctx) throws
AnalysisException {
+ // PROC headers belong to the resolved node and require the same
privilege as reading it.
+ if (!Env.getCurrentEnv().getAccessManager().checkGlobalPriv(ctx,
PrivPredicate.ADMIN_OR_NODE)) {
+
ErrorReport.reportAnalysisException(ErrorCode.ERR_SPECIFIC_ACCESS_DENIED_ERROR,
"ADMIN");
+ }
+ return getMetaData(ProcService.getInstance().open(path));
Review Comment:
[P2] Resolve SHOW PROC headers without fetching all rows. This new
schema-only method calls `getMetaData(procNode)`, which invokes
`fetchResult()`: `SHOW PROC '/current_queries'` builds and sorts every
active-query row, and remote index-schema nodes can contact BEs, all while
`FlightSqlQuerySchema.analyze` holds the session lock. Flight clients may call
Prepare before execution, repeating this work just to read column names. Use
header-only metadata where available, or decline schema discovery for nodes
that require a full fetch.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java:
##########
@@ -0,0 +1,301 @@
+// 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.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.Command;
+import org.apache.doris.nereids.trees.plans.commands.DescribeCommand;
+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.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.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());
+ statements = new NereidsParser().parseSQL(converted,
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 Prepare for ordinary no-row commands. `INSERT`, `UPDATE`,
`DELETE`, `MERGE INTO`, `KILL QUERY`, and `BEGIN` inherit empty command
metadata but their statement types miss this allowlist, so Prepare/GetSchema
now return UNIMPLEMENTED although direct Flight execution returns
`StatusResult`. Classify concrete command implementations or preserve the
prepared execution path: a blanket `INSERT` case is unsafe because `WARM UP
SELECT` can return six columns. Cover both result shapes in Flight tests.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java:
##########
@@ -0,0 +1,301 @@
+// 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.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.Command;
+import org.apache.doris.nereids.trees.plans.commands.DescribeCommand;
+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.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.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());
+ statements = new NereidsParser().parseSQL(converted,
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);
+ }
+ return command.getResultSetMetaData();
Review Comment:
[P2] Preserve preparation for the procedure SHOW commands. `SHOW PROCEDURE
STATUS` and `SHOW CREATE PROCEDURE` both return fixed `getMetaData()` columns
and execute through `sendResultSet`, but they extend `Command` without
overriding `getResultSetMetaData()`. This new SHOW fallback therefore returns
UNIMPLEMENTED during GetSchema and Prepare, whereas prepared execution was
previously accepted. Use their fixed headers (or add the metadata override)
before rejecting the command; these cases are separate from the earlier SHOW
CREATE TABLE and SHOW PROC threads.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java:
##########
@@ -0,0 +1,301 @@
+// 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.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.Command;
+import org.apache.doris.nereids.trees.plans.commands.DescribeCommand;
+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.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.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());
+ statements = new NereidsParser().parseSQL(converted,
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:
Review Comment:
[P2] Exclude result-producing ALTER variants from the OK schema. `CREATE
INDEX ... USING ANN` on a Lance table is an `AlterTableCommand`, so its empty
default metadata reaches this ALTER branch and Prepare/GetSchema advertise
`StatusResult`. `AlterTableCommand.run` instead calls `sendJobIdResult` and
Flight execution returns a `JobId` field, including for zero-row IF no-ops.
Resolve that variant's header or return UNIMPLEMENTED for it. The earlier
empty-OTHER thread does not cover ALTER.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQuerySchema.java:
##########
@@ -0,0 +1,301 @@
+// 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.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.Command;
+import org.apache.doris.nereids.trees.plans.commands.DescribeCommand;
+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.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.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());
+ statements = new NereidsParser().parseSQL(converted,
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 Prepare for `EXPLAIN SELECT 1`. `ExplainCommand` inherits
empty metadata and has `StmtType.EXPLAIN`, so this fallback now rejects its
GetSchema and Prepare calls; direct Flight execution instead returns the fixed
`Explain String(Nereids Planner)` column. `ReplayCommand` has the same gap and
returns `Plan Replayer dump url`. Resolve these headers without executing the
statements, accounting for EXPLAIN PLAN PROCESS's separate shape, or retain
their prior prepared path.
--
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]