This is an automated email from the ASF dual-hosted git repository.
dataroaring pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 1eaa1c9bee8 [fix](http stream) http stream should throw exception if
parse sql failed (#55863)
1eaa1c9bee8 is described below
commit 1eaa1c9bee836d53297d87c82164c1cd4688c79e
Author: meiyi <[email protected]>
AuthorDate: Thu Sep 11 09:04:38 2025 +0800
[fix](http stream) http stream should throw exception if parse sql failed
(#55863)
Problem Summary:
if http stream parse sql failed, the exception is ignored and continue
execute, then throw exception:
```
"Message": "[ANALYSIS_ERROR]TStatus: errCode = 2, detailMessage = exec
sql error catch unknown result.java.lang.NullPointerException: Cannot invoke
\"org.apache.doris.planner.Planner.getFragments()\" because \"planner\" is null"
```
---
.../main/java/org/apache/doris/qe/StmtExecutor.java | 12 +++++-------
.../org/apache/doris/service/FrontendServiceImpl.java | 4 ++--
.../load_p0/http_stream/test_http_stream.groovy | 19 +++++++++++++++++++
3 files changed, 26 insertions(+), 9 deletions(-)
diff --git a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
index 0f8bb745e9d..b994ea2bbbc 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/qe/StmtExecutor.java
@@ -1939,7 +1939,7 @@ public class StmtExecutor {
}
private HttpStreamParams generateHttpStreamNereidsPlan(TUniqueId queryId) {
- LOG.info("TUniqueId: {} generate stream load plan", queryId);
+ LOG.info("TUniqueId: {} generate stream load plan",
DebugUtil.printId(queryId));
context.setQueryId(queryId);
context.setStmtId(STMT_ID_GENERATOR.incrementAndGet());
@@ -1948,7 +1948,6 @@ public class StmtExecutor {
"Nereids only process LogicalPlanAdapter, but parsedStmt is "
+ parsedStmt.getClass().getName());
context.getState().setNereids(true);
InsertIntoTableCommand insert = (InsertIntoTableCommand)
((LogicalPlanAdapter) parsedStmt).getLogicalPlan();
- HttpStreamParams httpStreamParams = new HttpStreamParams();
try {
if
(!StringUtils.isEmpty(context.getSessionVariable().groupCommit)) {
@@ -1958,6 +1957,7 @@ public class StmtExecutor {
context.setGroupCommit(true);
}
OlapInsertExecutor insertExecutor = (OlapInsertExecutor)
insert.initPlan(context, this);
+ HttpStreamParams httpStreamParams = new HttpStreamParams();
httpStreamParams.setTxnId(insertExecutor.getTxnId());
httpStreamParams.setDb(insertExecutor.getDatabase());
httpStreamParams.setTable(insertExecutor.getTable());
@@ -1974,6 +1974,7 @@ public class StmtExecutor {
if (!isValidPlan) {
throw new AnalysisException("plan is invalid: " +
planRoot.getExplainString());
}
+ return httpStreamParams;
} catch (QueryStateException e) {
LOG.debug("Command(" + originStmt.originStmt + ") process
failed.", e);
context.setState(e.getQueryState());
@@ -1992,17 +1993,15 @@ public class StmtExecutor {
throw new NereidsException("Command (" + originStmt.originStmt +
") process failed.",
new AnalysisException(e.getMessage(), e));
}
- return httpStreamParams;
}
public HttpStreamParams generateHttpStreamPlan(TUniqueId queryId) throws
Exception {
SessionVariable sessionVariable = context.getSessionVariable();
- HttpStreamParams httpStreamParams = null;
try {
try {
// disable shuffle for http stream (only 1 sink)
sessionVariable.setVarOnce(SessionVariable.ENABLE_STRICT_CONSISTENCY_DML,
"false");
- httpStreamParams = generateHttpStreamNereidsPlan(queryId);
+ return generateHttpStreamNereidsPlan(queryId);
} catch (NereidsException | ParseException e) {
if (context.getMinidump() != null &&
context.getMinidump().toString(4) != null) {
MinidumpUtils.saveMinidumpString(context.getMinidump(),
DebugUtil.printId(context.queryId()));
@@ -2014,8 +2013,8 @@ public class StmtExecutor {
}
if (e instanceof NereidsException) {
LOG.warn("Analyze failed. {}",
context.getQueryIdentifier(), e);
- throw ((NereidsException) e).getException();
}
+ throw e;
} catch (Exception e) {
throw new RuntimeException(e);
}
@@ -2031,7 +2030,6 @@ public class StmtExecutor {
context.getState().setError(e.getMysqlErrorCode(),
e.getMessage());
}
}
- return httpStreamParams;
}
public SummaryProfile getSummaryProfile() {
diff --git
a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
index 83c53c532dc..21eb4a6c3ad 100644
--- a/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
+++ b/fe/fe-core/src/main/java/org/apache/doris/service/FrontendServiceImpl.java
@@ -2259,8 +2259,8 @@ public class FrontendServiceImpl implements
FrontendService.Iface {
LOG.warn("exec sql error", e);
throw e;
} catch (Throwable e) {
- LOG.warn("exec sql error catch unknown result.", e);
- throw new UserException("exec sql error catch unknown result." +
e);
+ LOG.warn("exec sql: {} catch unknown result. ", originStmt, e);
+ throw new UserException("exec sql error catch unknown result. " +
e.getMessage());
}
return httpStreamParams;
}
diff --git a/regression-test/suites/load_p0/http_stream/test_http_stream.groovy
b/regression-test/suites/load_p0/http_stream/test_http_stream.groovy
index f732f7ce3c2..afa9ca97252 100644
--- a/regression-test/suites/load_p0/http_stream/test_http_stream.groovy
+++ b/regression-test/suites/load_p0/http_stream/test_http_stream.groovy
@@ -63,6 +63,25 @@ suite("test_http_stream", "p0") {
}
}
+ // test error sql
+ streamLoad {
+ set 'version', '1'
+ set 'sql', """
+ insert into ${db}.${tableName1} (id, name) select
+ """
+ time 10000
+ file 'test_http_stream.csv'
+ check { result, exception, startTime, endTime ->
+ if (exception != null) {
+ throw exception
+ }
+ log.info("http_stream result: ${result}".toString())
+ def json = parseJson(result)
+ assertEquals("fail", json.Status.toLowerCase())
+ assertTrue(json.Message.contains("Nereids parse failed"))
+ }
+ }
+
qt_sql1 "select id, name from ${tableName1}"
} finally {
try_sql "DROP TABLE IF EXISTS ${tableName1}"
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]