This is an automated email from the ASF dual-hosted git repository.
jackietien pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new 257ed97251 [IOTDB-3303] Support last query api for new cluster (#6030)
257ed97251 is described below
commit 257ed9725176776eef9bca9bb9679084fcdb54a9
Author: Haonan <[email protected]>
AuthorDate: Fri May 27 18:29:21 2022 +0800
[IOTDB-3303] Support last query api for new cluster (#6030)
---
.../db/mpp/plan/parser/StatementGenerator.java | 22 ++++----
.../thrift/impl/DataNodeTSIServiceImpl.java | 58 +++++++++++++++++++++-
2 files changed, 65 insertions(+), 15 deletions(-)
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/StatementGenerator.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/StatementGenerator.java
index 7103246f38..41c19df05c 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/StatementGenerator.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/StatementGenerator.java
@@ -22,8 +22,6 @@ package org.apache.iotdb.db.mpp.plan.parser;
import org.apache.iotdb.commons.exception.IllegalPathException;
import org.apache.iotdb.commons.path.PartialPath;
import org.apache.iotdb.db.exception.query.QueryProcessException;
-import org.apache.iotdb.db.mpp.common.filter.BasicFunctionFilter;
-import org.apache.iotdb.db.mpp.plan.constant.FilterConstant;
import org.apache.iotdb.db.mpp.plan.expression.binary.GreaterEqualExpression;
import org.apache.iotdb.db.mpp.plan.expression.binary.LessThanExpression;
import org.apache.iotdb.db.mpp.plan.expression.binary.LogicAndExpression;
@@ -77,8 +75,6 @@ import java.time.ZoneId;
import java.util.ArrayList;
import java.util.List;
-import static org.apache.iotdb.commons.conf.IoTDBConstant.TIME;
-
/** Convert SQL and RPC requests to {@link Statement}. */
public class StatementGenerator {
@@ -100,7 +96,8 @@ public class StatementGenerator {
fromComponent.addPrefixPath(path);
}
selectComponent.addResultColumn(
- new ResultColumn(new TimeSeriesOperand(new PartialPath("")),
ResultColumn.ColumnType.RAW));
+ new ResultColumn(
+ new TimeSeriesOperand(new PartialPath("", false)),
ResultColumn.ColumnType.RAW));
// set query filter
GreaterEqualExpression leftPredicate =
@@ -136,16 +133,15 @@ public class StatementGenerator {
fromComponent.addPrefixPath(path);
}
selectComponent.addResultColumn(
- new ResultColumn(new TimeSeriesOperand(new PartialPath("")),
ResultColumn.ColumnType.RAW));
+ new ResultColumn(
+ new TimeSeriesOperand(new PartialPath("", false)),
ResultColumn.ColumnType.RAW));
// set query filter
- PartialPath timePath = new PartialPath(TIME, false);
- BasicFunctionFilter basicFunctionFilter =
- new BasicFunctionFilter(
- FilterConstant.FilterType.GREATERTHANOREQUALTO,
- timePath,
- Long.toString(lastDataQueryReq.getTime()));
- // whereCondition.setQueryFilter(basicFunctionFilter);
+ GreaterEqualExpression predicate =
+ new GreaterEqualExpression(
+ new TimestampOperand(),
+ new ConstantOperand(TSDataType.INT64,
Long.toString(lastDataQueryReq.getTime())));
+ whereCondition.setPredicate(predicate);
lastQueryStatement.setSelectComponent(selectComponent);
lastQueryStatement.setFromComponent(fromComponent);
diff --git
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
index 6f059e886f..7c3a5f6c9d 100644
---
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
+++
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeTSIServiceImpl.java
@@ -1026,7 +1026,7 @@ public class DataNodeTSIServiceImpl implements
TSIEventHandler {
}
QUERY_FREQUENCY_RECORDER.incrementAndGet();
- AUDIT_LOGGER.debug("Session {} execute Row Data Query: {}",
req.sessionId, req);
+ AUDIT_LOGGER.debug("Session {} execute Raw Data Query: {}",
req.sessionId, req);
long queryId = SESSION_MANAGER.requestQueryId(req.statementId, true);
// create and cache dataset
ExecutionResult result =
@@ -1070,7 +1070,61 @@ public class DataNodeTSIServiceImpl implements
TSIEventHandler {
@Override
public TSExecuteStatementResp executeLastDataQuery(TSLastDataQueryReq req) {
- throw new UnsupportedOperationException();
+ if (!SESSION_MANAGER.checkLogin(req.getSessionId())) {
+ return RpcUtils.getTSExecuteStatementResp(getNotLoggedInStatus());
+ }
+ long startTime = System.currentTimeMillis();
+ try {
+ Statement s =
+ StatementGenerator.createStatement(req,
SESSION_MANAGER.getZoneId(req.getSessionId()));
+
+ // permission check
+ TSStatus status = AuthorityChecker.checkAuthority(s, req.sessionId);
+ if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return RpcUtils.getTSExecuteStatementResp(status);
+ }
+
+ QUERY_FREQUENCY_RECORDER.incrementAndGet();
+ AUDIT_LOGGER.debug("Session {} execute Last Data Query: {}",
req.sessionId, req);
+ long queryId = SESSION_MANAGER.requestQueryId(req.statementId, true);
+ // create and cache dataset
+ ExecutionResult result =
+ COORDINATOR.execute(
+ s,
+ queryId,
+ SESSION_MANAGER.getSessionInfo(req.sessionId),
+ "",
+ PARTITION_FETCHER,
+ SCHEMA_FETCHER);
+
+ if (result.status.code != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ throw new RuntimeException("error code: " + result.status);
+ }
+
+ IQueryExecution queryExecution = COORDINATOR.getQueryExecution(queryId);
+
+ TSExecuteStatementResp resp;
+ if (queryExecution.isQuery()) {
+ resp = createResponse(queryExecution.getDatasetHeader(), queryId);
+ resp.setStatus(result.status);
+ resp.setQueryDataSet(
+ QueryDataSetUtils.convertTsBlockByFetchSize(queryExecution,
req.fetchSize));
+ } else {
+ resp = RpcUtils.getTSExecuteStatementResp(result.status);
+ }
+
+ return resp;
+ } catch (Exception e) {
+ // TODO call the coordinator to release query resource
+ return RpcUtils.getTSExecuteStatementResp(
+ onQueryException(e, "\"" + req + "\". " +
OperationType.EXECUTE_LAST_DATA_QUERY));
+ } finally {
+ addOperationLatency(Operation.EXECUTE_QUERY, startTime);
+ long costTime = System.currentTimeMillis() - startTime;
+ if (costTime >= CONFIG.getSlowQueryThreshold()) {
+ SLOW_SQL_LOGGER.info("Cost: {} ms, sql is {}", costTime, req);
+ }
+ }
}
@Override