This is an automated email from the ASF dual-hosted git repository.
zyk pushed a commit to branch rel/1.2
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rel/1.2 by this push:
new af5c45e3258 [To rel/1.2] Refactor Alter View (#10165)
af5c45e3258 is described below
commit af5c45e325807856432dbee81c48fabf1475d9c8
Author: Marcos_Zyk <[email protected]>
AuthorDate: Fri Jun 16 01:53:33 2023 +0800
[To rel/1.2] Refactor Alter View (#10165)
---
.../confignode/client/DataNodeRequestType.java | 2 +
.../client/async/AsyncDataNodeClientPool.java | 29 ++-
.../client/async/handlers/AsyncClientHandler.java | 5 +-
...RPCHandler.java => SchemaUpdateRPCHandler.java} | 6 +-
.../iotdb/confignode/manager/ConfigManager.java | 11 +
.../apache/iotdb/confignode/manager/IManager.java | 3 +
.../iotdb/confignode/manager/ProcedureManager.java | 52 +++-
...ocedure.java => AlterLogicalViewProcedure.java} | 288 ++++++++++-----------
.../impl/schema/DeleteLogicalViewProcedure.java | 3 +-
.../state/schema/AlterLogicalViewState.java | 25 ++
.../procedure/store/ProcedureFactory.java | 6 +
.../confignode/procedure/store/ProcedureType.java | 3 +-
.../thrift/ConfigNodeRPCServiceProcessor.java | 6 +
.../db/it/schema/view/IoTDBAliasSeriesIT.java | 13 +-
.../iotdb/db/it/schema/view/IoTDBAlterViewIT.java | 117 +++++++++
.../src/main/thrift/confignode.thrift | 7 +
.../thrift/src/main/thrift/datanode.thrift | 7 +
.../apache/iotdb/db/client/ConfigNodeClient.java | 22 ++
.../metadata/view/ViewNotExistException.java | 25 --
.../plan/schemaregion/SchemaRegionPlanType.java | 1 +
.../plan/schemaregion/SchemaRegionPlanVisitor.java | 5 +
.../impl/SchemaRegionPlanDeserializer.java | 9 +
.../impl/SchemaRegionPlanSerializer.java | 13 +
.../impl/SchemaRegionPlanTxtSerializer.java | 11 +
.../impl/write/AlterLogicalViewPlanImpl.java | 56 ++++
.../impl/write/SchemaRegionWritePlanFactory.java | 8 +
.../write/view/IAlterLogicalViewPlan.java | 46 ++++
.../db/metadata/schemaregion/ISchemaRegion.java | 3 +
.../schemaregion/SchemaRegionMemoryImpl.java | 23 ++
.../schemaregion/SchemaRegionSchemaFileImpl.java | 7 +
.../metadata/visitor/SchemaExecutionVisitor.java | 20 ++
.../config/executor/ClusterConfigTaskExecutor.java | 74 +++---
.../iotdb/db/mpp/plan/parser/ASTVisitor.java | 10 +-
.../db/mpp/plan/planner/LogicalPlanVisitor.java | 6 +-
.../mpp/plan/planner/plan/node/PlanNodeType.java | 6 +-
.../db/mpp/plan/planner/plan/node/PlanVisitor.java | 5 +
.../metedata/write/view/AlterLogicalViewNode.java | 186 +++++++++++++
.../metadata/view/AlterLogicalViewStatement.java | 11 -
.../metadata/view/CreateLogicalViewStatement.java | 4 +
.../impl/DataNodeInternalRPCServiceImpl.java | 35 +++
40 files changed, 920 insertions(+), 249 deletions(-)
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
index 80701acefa8..f4e839add37 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/DataNodeRequestType.java
@@ -95,6 +95,8 @@ public enum DataNodeRequestType {
ROLLBACK_VIEW_SCHEMA_BLACK_LIST,
DELETE_VIEW,
+ ALTER_VIEW,
+
/** @TODO Need to migrate to 'Node Maintenance' */
KILL_QUERY_INSTANCE,
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
index 018a46ad229..e410e095289 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/AsyncDataNodeClientPool.java
@@ -35,9 +35,10 @@ import
org.apache.iotdb.confignode.client.async.handlers.AsyncClientHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.AsyncTSStatusRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.CheckTimeSeriesExistenceRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.CountPathsUsingTemplateRPCHandler;
-import
org.apache.iotdb.confignode.client.async.handlers.rpc.DeleteSchemaRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.FetchSchemaBlackListRPCHandler;
+import
org.apache.iotdb.confignode.client.async.handlers.rpc.SchemaUpdateRPCHandler;
import org.apache.iotdb.mpp.rpc.thrift.TActiveTriggerInstanceReq;
+import org.apache.iotdb.mpp.rpc.thrift.TAlterViewReq;
import org.apache.iotdb.mpp.rpc.thrift.TCheckTimeSeriesExistenceReq;
import org.apache.iotdb.mpp.rpc.thrift.TConstructSchemaBlackListReq;
import
org.apache.iotdb.mpp.rpc.thrift.TConstructSchemaBlackListWithTemplateReq;
@@ -278,13 +279,13 @@ public class AsyncDataNodeClientPool {
case CONSTRUCT_SCHEMA_BLACK_LIST:
client.constructSchemaBlackList(
(TConstructSchemaBlackListReq)
clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case ROLLBACK_SCHEMA_BLACK_LIST:
client.rollbackSchemaBlackList(
(TRollbackSchemaBlackListReq)
clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case FETCH_SCHEMA_BLACK_LIST:
@@ -302,31 +303,31 @@ public class AsyncDataNodeClientPool {
case DELETE_DATA_FOR_DELETE_SCHEMA:
client.deleteDataForDeleteSchema(
(TDeleteDataForDeleteSchemaReq)
clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case DELETE_TIMESERIES:
client.deleteTimeSeries(
(TDeleteTimeSeriesReq) clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case CONSTRUCT_SCHEMA_BLACK_LIST_WITH_TEMPLATE:
client.constructSchemaBlackListWithTemplate(
(TConstructSchemaBlackListWithTemplateReq)
clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case ROLLBACK_SCHEMA_BLACK_LIST_WITH_TEMPLATE:
client.rollbackSchemaBlackListWithTemplate(
(TRollbackSchemaBlackListWithTemplateReq)
clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case DEACTIVATE_TEMPLATE:
client.deactivateTemplate(
(TDeactivateTemplateReq) clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case UPDATE_TEMPLATE:
@@ -350,19 +351,25 @@ public class AsyncDataNodeClientPool {
case CONSTRUCT_VIEW_SCHEMA_BLACK_LIST:
client.constructViewSchemaBlackList(
(TConstructViewSchemaBlackListReq)
clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case ROLLBACK_VIEW_SCHEMA_BLACK_LIST:
client.rollbackViewSchemaBlackList(
(TRollbackViewSchemaBlackListReq)
clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case DELETE_VIEW:
client.deleteViewSchema(
(TDeleteViewSchemaReq) clientHandler.getRequest(requestId),
- (DeleteSchemaRPCHandler)
+ (SchemaUpdateRPCHandler)
+ clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
+ break;
+ case ALTER_VIEW:
+ client.alterView(
+ (TAlterViewReq) clientHandler.getRequest(requestId),
+ (SchemaUpdateRPCHandler)
clientHandler.createAsyncRPCHandler(requestId,
targetDataNode));
break;
case KILL_QUERY_INSTANCE:
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
index a37c225a04a..efd4b9f8f55 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/AsyncClientHandler.java
@@ -25,8 +25,8 @@ import
org.apache.iotdb.confignode.client.async.handlers.rpc.AbstractAsyncRPCHan
import
org.apache.iotdb.confignode.client.async.handlers.rpc.AsyncTSStatusRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.CheckTimeSeriesExistenceRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.CountPathsUsingTemplateRPCHandler;
-import
org.apache.iotdb.confignode.client.async.handlers.rpc.DeleteSchemaRPCHandler;
import
org.apache.iotdb.confignode.client.async.handlers.rpc.FetchSchemaBlackListRPCHandler;
+import
org.apache.iotdb.confignode.client.async.handlers.rpc.SchemaUpdateRPCHandler;
import org.apache.iotdb.mpp.rpc.thrift.TCheckTimeSeriesExistenceResp;
import org.apache.iotdb.mpp.rpc.thrift.TCountPathsUsingTemplateResp;
import org.apache.iotdb.mpp.rpc.thrift.TFetchSchemaBlackListResp;
@@ -165,7 +165,8 @@ public class AsyncClientHandler<Q, R> {
case CONSTRUCT_VIEW_SCHEMA_BLACK_LIST:
case ROLLBACK_VIEW_SCHEMA_BLACK_LIST:
case DELETE_VIEW:
- return new DeleteSchemaRPCHandler(
+ case ALTER_VIEW:
+ return new SchemaUpdateRPCHandler(
requestType,
requestId,
targetDataNode,
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DeleteSchemaRPCHandler.java
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SchemaUpdateRPCHandler.java
similarity index 93%
rename from
confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DeleteSchemaRPCHandler.java
rename to
confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SchemaUpdateRPCHandler.java
index 0b7fab3f407..9043a653cdf 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/DeleteSchemaRPCHandler.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/client/async/handlers/rpc/SchemaUpdateRPCHandler.java
@@ -31,11 +31,11 @@ import org.slf4j.LoggerFactory;
import java.util.Map;
import java.util.concurrent.CountDownLatch;
-public class DeleteSchemaRPCHandler extends AsyncTSStatusRPCHandler {
+public class SchemaUpdateRPCHandler extends AsyncTSStatusRPCHandler {
- private static final Logger LOGGER =
LoggerFactory.getLogger(DeleteSchemaRPCHandler.class);
+ private static final Logger LOGGER =
LoggerFactory.getLogger(SchemaUpdateRPCHandler.class);
- public DeleteSchemaRPCHandler(
+ public SchemaUpdateRPCHandler(
DataNodeRequestType requestType,
int requestId,
TDataNodeLocation targetDataNode,
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
index 42baef81f85..18880cf9597 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ConfigManager.java
@@ -101,6 +101,7 @@ import
org.apache.iotdb.confignode.persistence.partition.PartitionInfo;
import org.apache.iotdb.confignode.persistence.pipe.PipeInfo;
import org.apache.iotdb.confignode.persistence.quota.QuotaInfo;
import org.apache.iotdb.confignode.persistence.schema.ClusterSchemaInfo;
+import org.apache.iotdb.confignode.rpc.thrift.TAlterLogicalViewReq;
import org.apache.iotdb.confignode.rpc.thrift.TAlterSchemaTemplateReq;
import org.apache.iotdb.confignode.rpc.thrift.TClusterParameters;
import org.apache.iotdb.confignode.rpc.thrift.TConfigNodeRegisterReq;
@@ -1593,6 +1594,16 @@ public class ConfigManager implements IManager {
}
}
+ @Override
+ public TSStatus alterLogicalView(TAlterLogicalViewReq req) {
+ TSStatus status = confirmLeader();
+ if (status.getCode() == TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
+ return procedureManager.alterLogicalView(req);
+ } else {
+ return status;
+ }
+ }
+
@Override
public TSStatus createPipe(TCreatePipeReq req) {
TSStatus status = confirmLeader();
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/IManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/IManager.java
index 9e8560d0aec..0a40378911b 100644
--- a/confignode/src/main/java/org/apache/iotdb/confignode/manager/IManager.java
+++ b/confignode/src/main/java/org/apache/iotdb/confignode/manager/IManager.java
@@ -46,6 +46,7 @@ import org.apache.iotdb.confignode.manager.node.NodeManager;
import org.apache.iotdb.confignode.manager.partition.PartitionManager;
import org.apache.iotdb.confignode.manager.pipe.PipeManager;
import org.apache.iotdb.confignode.manager.schema.ClusterSchemaManager;
+import org.apache.iotdb.confignode.rpc.thrift.TAlterLogicalViewReq;
import org.apache.iotdb.confignode.rpc.thrift.TAlterSchemaTemplateReq;
import org.apache.iotdb.confignode.rpc.thrift.TConfigNodeRegisterReq;
import org.apache.iotdb.confignode.rpc.thrift.TConfigNodeRegisterResp;
@@ -553,6 +554,8 @@ public interface IManager {
TSStatus deleteLogicalView(TDeleteLogicalViewReq req);
+ TSStatus alterLogicalView(TAlterLogicalViewReq req);
+
/**
* Create Pipe
*
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
index d842f8dc3ad..fc84b840245 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/manager/ProcedureManager.java
@@ -29,8 +29,10 @@ import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.exception.IoTDBException;
import org.apache.iotdb.commons.model.ModelInformation;
import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.path.PathDeserializeUtil;
import org.apache.iotdb.commons.path.PathPatternTree;
import org.apache.iotdb.commons.pipe.plugin.meta.PipePluginMeta;
+import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
import org.apache.iotdb.commons.trigger.TriggerInformation;
import org.apache.iotdb.commons.utils.StatusUtils;
import org.apache.iotdb.confignode.conf.ConfigNodeConfig;
@@ -59,6 +61,7 @@ import
org.apache.iotdb.confignode.procedure.impl.pipe.task.CreatePipeProcedureV
import
org.apache.iotdb.confignode.procedure.impl.pipe.task.DropPipeProcedureV2;
import
org.apache.iotdb.confignode.procedure.impl.pipe.task.StartPipeProcedureV2;
import
org.apache.iotdb.confignode.procedure.impl.pipe.task.StopPipeProcedureV2;
+import
org.apache.iotdb.confignode.procedure.impl.schema.AlterLogicalViewProcedure;
import
org.apache.iotdb.confignode.procedure.impl.schema.DeactivateTemplateProcedure;
import
org.apache.iotdb.confignode.procedure.impl.schema.DeleteDatabaseProcedure;
import
org.apache.iotdb.confignode.procedure.impl.schema.DeleteLogicalViewProcedure;
@@ -76,6 +79,7 @@ import
org.apache.iotdb.confignode.procedure.store.IProcedureStore;
import org.apache.iotdb.confignode.procedure.store.ProcedureFactory;
import org.apache.iotdb.confignode.procedure.store.ProcedureStore;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
+import org.apache.iotdb.confignode.rpc.thrift.TAlterLogicalViewReq;
import org.apache.iotdb.confignode.rpc.thrift.TConfigNodeRegisterReq;
import org.apache.iotdb.confignode.rpc.thrift.TCreateCQReq;
import org.apache.iotdb.confignode.rpc.thrift.TCreatePipeReq;
@@ -98,6 +102,7 @@ import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Set;
@@ -234,7 +239,7 @@ public class ProcedureManager {
DeleteLogicalViewProcedure deleteLogicalViewProcedure;
for (Procedure<?> procedure : executor.getProcedures().values()) {
type = ProcedureFactory.getProcedureType(procedure);
- if (type == null ||
!type.equals(ProcedureType.DELETE_TIMESERIES_PROCEDURE)) {
+ if (type == null ||
!type.equals(ProcedureType.DELETE_LOGICAL_VIEW_PROCEDURE)) {
continue;
}
deleteLogicalViewProcedure = ((DeleteLogicalViewProcedure) procedure);
@@ -268,6 +273,51 @@ public class ProcedureManager {
}
}
+ public TSStatus alterLogicalView(TAlterLogicalViewReq req) {
+ String queryId = req.getQueryId();
+ ByteBuffer byteBuffer = ByteBuffer.wrap(req.getViewBinary());
+ Map<PartialPath, ViewExpression> viewPathToSourceMap = new HashMap<>();
+ int size = byteBuffer.getInt();
+ PartialPath path;
+ ViewExpression viewExpression;
+ for (int i = 0; i < size; i++) {
+ path = (PartialPath) PathDeserializeUtil.deserialize(byteBuffer);
+ viewExpression = ViewExpression.deserialize(byteBuffer);
+ viewPathToSourceMap.put(path, viewExpression);
+ }
+
+ long procedureId = -1;
+ synchronized (this) {
+ ProcedureType type;
+ AlterLogicalViewProcedure alterLogicalViewProcedure;
+ for (Procedure<?> procedure : executor.getProcedures().values()) {
+ type = ProcedureFactory.getProcedureType(procedure);
+ if (type == null ||
!type.equals(ProcedureType.ALTER_LOGICAL_VIEW_PROCEDURE)) {
+ continue;
+ }
+ alterLogicalViewProcedure = ((AlterLogicalViewProcedure) procedure);
+ if (queryId.equals(alterLogicalViewProcedure.getQueryId())) {
+ procedureId = alterLogicalViewProcedure.getProcId();
+ break;
+ }
+ }
+
+ if (procedureId == -1) {
+ procedureId =
+ this.executor.submitProcedure(
+ new AlterLogicalViewProcedure(queryId, viewPathToSourceMap));
+ }
+ }
+ List<TSStatus> procedureStatus = new ArrayList<>();
+ boolean isSucceed =
+ waitingProcedureFinished(Collections.singletonList(procedureId),
procedureStatus);
+ if (isSucceed) {
+ return StatusUtils.OK;
+ } else {
+ return procedureStatus.get(0);
+ }
+ }
+
public TSStatus setSchemaTemplate(String queryId, String templateName,
String templateSetPath) {
long procedureId = -1;
synchronized (this) {
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/AlterLogicalViewProcedure.java
similarity index 51%
copy from
confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
copy to
confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/AlterLogicalViewProcedure.java
index 8547eeaed6a..c1802db6214 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/AlterLogicalViewProcedure.java
@@ -23,9 +23,12 @@ import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId;
import org.apache.iotdb.common.rpc.thrift.TDataNodeLocation;
import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
import org.apache.iotdb.common.rpc.thrift.TSStatus;
+import org.apache.iotdb.common.rpc.thrift.TSeriesPartitionSlot;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.path.PathDeserializeUtil;
import org.apache.iotdb.commons.path.PathPatternTree;
+import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
import org.apache.iotdb.confignode.client.DataNodeRequestType;
import org.apache.iotdb.confignode.client.async.AsyncDataNodeClientPool;
import org.apache.iotdb.confignode.client.async.handlers.AsyncClientHandler;
@@ -34,13 +37,11 @@ import
org.apache.iotdb.confignode.procedure.exception.ProcedureException;
import
org.apache.iotdb.confignode.procedure.exception.ProcedureSuspendedException;
import org.apache.iotdb.confignode.procedure.exception.ProcedureYieldException;
import
org.apache.iotdb.confignode.procedure.impl.statemachine.StateMachineProcedure;
-import
org.apache.iotdb.confignode.procedure.state.schema.DeleteLogicalViewState;
+import
org.apache.iotdb.confignode.procedure.state.schema.AlterLogicalViewState;
import org.apache.iotdb.confignode.procedure.store.ProcedureType;
import org.apache.iotdb.db.exception.metadata.view.ViewNotExistException;
-import org.apache.iotdb.mpp.rpc.thrift.TConstructViewSchemaBlackListReq;
-import org.apache.iotdb.mpp.rpc.thrift.TDeleteViewSchemaReq;
+import org.apache.iotdb.mpp.rpc.thrift.TAlterViewReq;
import org.apache.iotdb.mpp.rpc.thrift.TInvalidateMatchedSchemaCacheReq;
-import org.apache.iotdb.mpp.rpc.thrift.TRollbackViewSchemaBlackListReq;
import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
@@ -52,129 +53,68 @@ import java.io.DataOutputStream;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.function.BiFunction;
-import java.util.stream.Collectors;
-public class DeleteLogicalViewProcedure
- extends StateMachineProcedure<ConfigNodeProcedureEnv,
DeleteLogicalViewState> {
+public class AlterLogicalViewProcedure
+ extends StateMachineProcedure<ConfigNodeProcedureEnv,
AlterLogicalViewState> {
- private static final Logger LOGGER =
LoggerFactory.getLogger(DeleteLogicalViewProcedure.class);
+ private static final Logger LOGGER =
LoggerFactory.getLogger(AlterLogicalViewProcedure.class);
private String queryId;
- private PathPatternTree patternTree;
- private transient ByteBuffer patternTreeBytes;
+ private Map<PartialPath, ViewExpression> viewPathToSourceMap;
- private transient String requestMessage;
+ private transient PathPatternTree pathPatternTree;
+ private transient ByteBuffer patternTreeBytes;
- public DeleteLogicalViewProcedure() {
+ public AlterLogicalViewProcedure() {
super();
}
- public DeleteLogicalViewProcedure(String queryId, PathPatternTree
patternTree) {
+ public AlterLogicalViewProcedure(
+ String queryId, Map<PartialPath, ViewExpression> viewPathToSourceMap) {
super();
this.queryId = queryId;
- setPatternTree(patternTree);
+ this.viewPathToSourceMap = viewPathToSourceMap;
+ generatePathPatternTree();
}
@Override
- protected Flow executeFromState(ConfigNodeProcedureEnv env,
DeleteLogicalViewState state)
+ protected Flow executeFromState(ConfigNodeProcedureEnv env,
AlterLogicalViewState state)
throws ProcedureSuspendedException, ProcedureYieldException,
InterruptedException {
long startTime = System.currentTimeMillis();
try {
switch (state) {
- case CONSTRUCT_BLACK_LIST:
- LOGGER.info("Construct view schema black list of view {}",
requestMessage);
- if (constructBlackList(env) > 0) {
- setNextState(DeleteLogicalViewState.CLEAN_DATANODE_SCHEMA_CACHE);
- break;
- } else {
- setFailure(
- new ProcedureException(
- new ViewNotExistException(
- patternTree.getAllPathPatterns().stream()
- .map(PartialPath::getFullPath)
- .collect(Collectors.toList()),
- false)));
- return Flow.NO_MORE_STATE;
- }
case CLEAN_DATANODE_SCHEMA_CACHE:
- LOGGER.info("Invalidate cache of view {}", requestMessage);
+ LOGGER.info("Invalidate cache of view {}",
viewPathToSourceMap.keySet());
invalidateCache(env);
- break;
- case DELETE_VIEW_SCHEMA:
- LOGGER.info("Delete view schema of {}", requestMessage);
- deleteViewSchema(env);
+ setNextState(AlterLogicalViewState.ALTER_LOGICAL_VIEW);
+ return Flow.HAS_MORE_STATE;
+ case ALTER_LOGICAL_VIEW:
+ LOGGER.info("Alter view {}", viewPathToSourceMap.keySet());
+ try {
+ alterLogicalView(env);
+ } catch (ProcedureException e) {
+ setFailure(e);
+ }
return Flow.NO_MORE_STATE;
default:
setFailure(new ProcedureException("Unrecognized state " +
state.toString()));
return Flow.NO_MORE_STATE;
}
- return Flow.HAS_MORE_STATE;
} finally {
LOGGER.info(
String.format(
- "DeleteLogicalView-[%s] costs %sms",
+ "AlterLogicalView-[%s] costs %sms",
state.toString(), (System.currentTimeMillis() - startTime)));
}
}
- // return the total num of timeseries in schema black list
- private long constructBlackList(ConfigNodeProcedureEnv env) {
- Map<TConsensusGroupId, TRegionReplicaSet> targetSchemaRegionGroup =
- env.getConfigManager().getRelatedSchemaRegionGroup(patternTree);
- if (targetSchemaRegionGroup.isEmpty()) {
- return 0;
- }
- List<TSStatus> successResult = new ArrayList<>();
- DeleteLogicalViewRegionTaskExecutor<TConstructViewSchemaBlackListReq>
constructBlackListTask =
- new
DeleteLogicalViewRegionTaskExecutor<TConstructViewSchemaBlackListReq>(
- "construct view schema black list",
- env,
- targetSchemaRegionGroup,
- DataNodeRequestType.CONSTRUCT_VIEW_SCHEMA_BLACK_LIST,
- ((dataNodeLocation, consensusGroupIdList) ->
- new TConstructViewSchemaBlackListReq(consensusGroupIdList,
patternTreeBytes))) {
- @Override
- protected List<TConsensusGroupId> processResponseOfOneDataNode(
- TDataNodeLocation dataNodeLocation,
- List<TConsensusGroupId> consensusGroupIdList,
- TSStatus response) {
- List<TConsensusGroupId> failedRegionList = new ArrayList<>();
- if (response.getCode() ==
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- successResult.add(response);
- } else if (response.getCode() ==
TSStatusCode.MULTIPLE_ERROR.getStatusCode()) {
- List<TSStatus> subStatusList = response.getSubStatus();
- for (int i = 0; i < subStatusList.size(); i++) {
- if (subStatusList.get(i).getCode() ==
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- successResult.add(subStatusList.get(i));
- } else {
- failedRegionList.add(consensusGroupIdList.get(i));
- }
- }
- } else {
- failedRegionList.addAll(consensusGroupIdList);
- }
- return failedRegionList;
- }
- };
- constructBlackListTask.execute();
-
- if (isFailed()) {
- return 0;
- }
-
- long preDeletedNum = 0;
- for (TSStatus resp : successResult) {
- preDeletedNum += Long.parseLong(resp.getMessage());
- }
- return preDeletedNum;
- }
-
private void invalidateCache(ConfigNodeProcedureEnv env) {
Map<Integer, TDataNodeLocation> dataNodeLocationMap =
env.getConfigManager().getNodeManager().getRegisteredDataNodeLocations();
@@ -188,78 +128,114 @@ public class DeleteLogicalViewProcedure
for (TSStatus status : statusMap.values()) {
// all dataNodes must clear the related schema cache
if (status.getCode() != TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- LOGGER.error("Failed to invalidate schema cache of view {}",
requestMessage);
+ LOGGER.error("Failed to invalidate schema cache of view {}",
viewPathToSourceMap.keySet());
setFailure(
new ProcedureException(new MetadataException("Invalidate view
schema cache failed")));
return;
}
}
-
- setNextState(DeleteLogicalViewState.DELETE_VIEW_SCHEMA);
}
- private void deleteViewSchema(ConfigNodeProcedureEnv env) {
- DeleteLogicalViewRegionTaskExecutor<TDeleteViewSchemaReq>
deleteTimeSeriesTask =
- new DeleteLogicalViewRegionTaskExecutor<>(
- "delete view schema",
+ private void alterLogicalView(ConfigNodeProcedureEnv env) throws
ProcedureException {
+ Map<TConsensusGroupId, TRegionReplicaSet> targetSchemaRegionGroup =
+ env.getConfigManager().getRelatedSchemaRegionGroup(pathPatternTree);
+ Map<TConsensusGroupId, Map<PartialPath, ViewExpression>>
schemaRegionRequestMap =
+ new HashMap<>();
+ for (Map.Entry<PartialPath, ViewExpression> entry :
viewPathToSourceMap.entrySet()) {
+ schemaRegionRequestMap
+ .computeIfAbsent(getBelongedSchemaRegion(env, entry.getKey()), k ->
new HashMap<>())
+ .put(entry.getKey(), entry.getValue());
+ }
+ AlterLogicalViewRegionTaskExecutor<TAlterViewReq> regionTaskExecutor =
+ new AlterLogicalViewRegionTaskExecutor<>(
+ "Alter view",
env,
- env.getConfigManager().getRelatedSchemaRegionGroup(patternTree),
- DataNodeRequestType.DELETE_VIEW,
- ((dataNodeLocation, consensusGroupIdList) ->
- new TDeleteViewSchemaReq(consensusGroupIdList,
patternTreeBytes)));
- deleteTimeSeriesTask.execute();
+ targetSchemaRegionGroup,
+ DataNodeRequestType.ALTER_VIEW,
+ (dataNodeLocation, consensusGroupIdList) -> {
+ TAlterViewReq req = new TAlterViewReq();
+ req.setSchemaRegionIdList(consensusGroupIdList);
+ List<ByteBuffer> viewMapBinaryList = new ArrayList<>();
+ for (TConsensusGroupId consensusGroupId : consensusGroupIdList) {
+ ByteArrayOutputStream stream = new ByteArrayOutputStream();
+ Map<PartialPath, ViewExpression> viewMap =
+ schemaRegionRequestMap.get(consensusGroupId);
+ try {
+ ReadWriteIOUtils.write(viewMap.size(), stream);
+ for (Map.Entry<PartialPath, ViewExpression> viewEntry :
viewMap.entrySet()) {
+ viewEntry.getKey().serialize(stream);
+ ViewExpression.serialize(viewEntry.getValue(), stream);
+ }
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+ viewMapBinaryList.add(ByteBuffer.wrap(stream.toByteArray()));
+ }
+ req.setViewBinaryList(viewMapBinaryList);
+ return req;
+ });
+ regionTaskExecutor.execute();
+
+ invalidateCache(env);
}
- @Override
- protected void rollbackState(
- ConfigNodeProcedureEnv env, DeleteLogicalViewState
deleteLogicalViewState)
- throws IOException, InterruptedException, ProcedureException {
- DeleteLogicalViewRegionTaskExecutor<TRollbackViewSchemaBlackListReq>
rollbackStateTask =
- new DeleteLogicalViewRegionTaskExecutor<>(
- "roll back view schema black list",
- env,
- env.getConfigManager().getRelatedSchemaRegionGroup(patternTree),
- DataNodeRequestType.ROLLBACK_VIEW_SCHEMA_BLACK_LIST,
- (dataNodeLocation, consensusGroupIdList) ->
- new TRollbackViewSchemaBlackListReq(consensusGroupIdList,
patternTreeBytes));
- rollbackStateTask.execute();
+ private TConsensusGroupId getBelongedSchemaRegion(
+ ConfigNodeProcedureEnv env, PartialPath viewPath) throws
ProcedureException {
+ PathPatternTree patternTree = new PathPatternTree();
+ patternTree.appendFullPath(viewPath);
+ patternTree.constructTree();
+ Map<String, Map<TSeriesPartitionSlot, TConsensusGroupId>>
schemaPartitionTable =
+
env.getConfigManager().getSchemaPartition(patternTree).schemaPartitionTable;
+ if (schemaPartitionTable.isEmpty()) {
+ throw new ProcedureException(new
ViewNotExistException(viewPath.getFullPath()));
+ } else {
+ Map<TSeriesPartitionSlot, TConsensusGroupId> slotMap =
+ schemaPartitionTable.values().iterator().next();
+ if (slotMap.isEmpty()) {
+ throw new ProcedureException(new
ViewNotExistException(viewPath.getFullPath()));
+ } else {
+ return slotMap.values().iterator().next();
+ }
+ }
}
@Override
- protected boolean isRollbackSupported(DeleteLogicalViewState
deleteLogicalViewState) {
+ protected boolean isRollbackSupported(AlterLogicalViewState
alterLogicalViewState) {
return true;
}
@Override
- protected DeleteLogicalViewState getState(int stateId) {
- return DeleteLogicalViewState.values()[stateId];
+ protected void rollbackState(
+ ConfigNodeProcedureEnv env, AlterLogicalViewState alterLogicalViewState)
+ throws IOException, InterruptedException, ProcedureException {
+ invalidateCache(env);
}
@Override
- protected int getStateId(DeleteLogicalViewState deleteLogicalViewState) {
- return deleteLogicalViewState.ordinal();
+ protected AlterLogicalViewState getState(int stateId) {
+ return AlterLogicalViewState.values()[stateId];
}
@Override
- protected DeleteLogicalViewState getInitialState() {
- return DeleteLogicalViewState.CONSTRUCT_BLACK_LIST;
- }
-
- public String getQueryId() {
- return queryId;
+ protected int getStateId(AlterLogicalViewState alterLogicalViewState) {
+ return alterLogicalViewState.ordinal();
}
- public PathPatternTree getPatternTree() {
- return patternTree;
+ @Override
+ protected AlterLogicalViewState getInitialState() {
+ return AlterLogicalViewState.CLEAN_DATANODE_SCHEMA_CACHE;
}
- public void setPatternTree(PathPatternTree patternTree) {
- this.patternTree = patternTree;
- requestMessage = patternTree.getAllPathPatterns().toString();
- patternTreeBytes = preparePatternTreeBytesData(patternTree);
+ public String getQueryId() {
+ return queryId;
}
- private ByteBuffer preparePatternTreeBytesData(PathPatternTree patternTree) {
+ private void generatePathPatternTree() {
+ PathPatternTree patternTree = new PathPatternTree();
+ for (PartialPath path : viewPathToSourceMap.keySet()) {
+ patternTree.appendFullPath(path);
+ }
+ patternTree.constructTree();
ByteArrayOutputStream byteArrayOutputStream = new ByteArrayOutputStream();
DataOutputStream dataOutputStream = new
DataOutputStream(byteArrayOutputStream);
try {
@@ -267,45 +243,62 @@ public class DeleteLogicalViewProcedure
} catch (IOException ignored) {
}
- return ByteBuffer.wrap(byteArrayOutputStream.toByteArray());
+ ByteBuffer patternTreeBytes =
ByteBuffer.wrap(byteArrayOutputStream.toByteArray());
+
+ this.pathPatternTree = patternTree;
+ this.patternTreeBytes = patternTreeBytes;
}
@Override
public void serialize(DataOutputStream stream) throws IOException {
-
stream.writeShort(ProcedureType.DELETE_LOGICAL_VIEW_PROCEDURE.getTypeCode());
+ stream.writeInt(ProcedureType.ALTER_LOGICAL_VIEW_PROCEDURE.getTypeCode());
super.serialize(stream);
ReadWriteIOUtils.write(queryId, stream);
- patternTree.serialize(stream);
+ ReadWriteIOUtils.write(this.viewPathToSourceMap.size(), stream);
+ for (Map.Entry<PartialPath, ViewExpression> entry :
viewPathToSourceMap.entrySet()) {
+ entry.getKey().serialize(stream);
+ ViewExpression.serialize(entry.getValue(), stream);
+ }
}
@Override
public void deserialize(ByteBuffer byteBuffer) {
super.deserialize(byteBuffer);
queryId = ReadWriteIOUtils.readString(byteBuffer);
- setPatternTree(PathPatternTree.deserialize(byteBuffer));
+
+ Map<PartialPath, ViewExpression> viewPathToSourceMap = new HashMap<>();
+ int size = byteBuffer.getInt();
+ PartialPath path;
+ ViewExpression viewExpression;
+ for (int i = 0; i < size; i++) {
+ path = (PartialPath) PathDeserializeUtil.deserialize(byteBuffer);
+ viewExpression = ViewExpression.deserialize(byteBuffer);
+ viewPathToSourceMap.put(path, viewExpression);
+ }
+ this.viewPathToSourceMap = viewPathToSourceMap;
+ generatePathPatternTree();
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
- if (o == null || getClass() != o.getClass()) return false;
- DeleteLogicalViewProcedure that = (DeleteLogicalViewProcedure) o;
- return this.getProcId() == that.getProcId()
- && this.getState() == that.getState()
- && patternTree.equals(that.patternTree);
+ if (!(o instanceof AlterLogicalViewProcedure)) return false;
+ AlterLogicalViewProcedure that = (AlterLogicalViewProcedure) o;
+ return Objects.equals(queryId, that.queryId)
+ && Objects.equals(viewPathToSourceMap, that.viewPathToSourceMap);
}
@Override
public int hashCode() {
- return Objects.hash(getProcId(), getState(), patternTree);
+ return Objects.hash(queryId, viewPathToSourceMap);
}
- private class DeleteLogicalViewRegionTaskExecutor<Q>
+ private class AlterLogicalViewRegionTaskExecutor<Q>
extends DataNodeRegionTaskExecutor<Q, TSStatus> {
private final String taskName;
- DeleteLogicalViewRegionTaskExecutor(
+ AlterLogicalViewRegionTaskExecutor(
String taskName,
ConfigNodeProcedureEnv env,
Map<TConsensusGroupId, TRegionReplicaSet> targetSchemaRegionGroup,
@@ -345,8 +338,11 @@ public class DeleteLogicalViewProcedure
new ProcedureException(
new MetadataException(
String.format(
- "Delete view %s failed when [%s] because all replicaset
of schemaRegion %s failed. %s",
- requestMessage, taskName, consensusGroupId.id,
dataNodeLocationSet))));
+ "Alter view %s failed when [%s] because all replicaset
of schemaRegion %s failed. %s",
+ viewPathToSourceMap.keySet(),
+ taskName,
+ consensusGroupId.id,
+ dataNodeLocationSet))));
interruptTask();
}
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
index 8547eeaed6a..4510401cdbb 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/impl/schema/DeleteLogicalViewProcedure.java
@@ -98,8 +98,7 @@ public class DeleteLogicalViewProcedure
new ViewNotExistException(
patternTree.getAllPathPatterns().stream()
.map(PartialPath::getFullPath)
- .collect(Collectors.toList()),
- false)));
+ .collect(Collectors.toList()))));
return Flow.NO_MORE_STATE;
}
case CLEAN_DATANODE_SCHEMA_CACHE:
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/schema/AlterLogicalViewState.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/schema/AlterLogicalViewState.java
new file mode 100644
index 00000000000..4885f981769
--- /dev/null
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/state/schema/AlterLogicalViewState.java
@@ -0,0 +1,25 @@
+/*
+ * 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.iotdb.confignode.procedure.state.schema;
+
+public enum AlterLogicalViewState {
+ CLEAN_DATANODE_SCHEMA_CACHE,
+ ALTER_LOGICAL_VIEW
+}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java
index e352f606312..67d8583b58f 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureFactory.java
@@ -35,6 +35,7 @@ import
org.apache.iotdb.confignode.procedure.impl.pipe.task.CreatePipeProcedureV
import
org.apache.iotdb.confignode.procedure.impl.pipe.task.DropPipeProcedureV2;
import
org.apache.iotdb.confignode.procedure.impl.pipe.task.StartPipeProcedureV2;
import
org.apache.iotdb.confignode.procedure.impl.pipe.task.StopPipeProcedureV2;
+import
org.apache.iotdb.confignode.procedure.impl.schema.AlterLogicalViewProcedure;
import
org.apache.iotdb.confignode.procedure.impl.schema.DeactivateTemplateProcedure;
import
org.apache.iotdb.confignode.procedure.impl.schema.DeleteDatabaseProcedure;
import
org.apache.iotdb.confignode.procedure.impl.schema.DeleteLogicalViewProcedure;
@@ -96,6 +97,9 @@ public class ProcedureFactory implements IProcedureFactory {
case DELETE_LOGICAL_VIEW_PROCEDURE:
procedure = new DeleteLogicalViewProcedure();
break;
+ case ALTER_LOGICAL_VIEW_PROCEDURE:
+ procedure = new AlterLogicalViewProcedure();
+ break;
case CREATE_TRIGGER_PROCEDURE:
procedure = new CreateTriggerProcedure();
break;
@@ -228,6 +232,8 @@ public class ProcedureFactory implements IProcedureFactory {
return ProcedureType.PIPE_HANDLE_META_CHANGE_PROCEDURE;
} else if (procedure instanceof DeleteLogicalViewProcedure) {
return ProcedureType.DELETE_LOGICAL_VIEW_PROCEDURE;
+ } else if (procedure instanceof AlterLogicalViewProcedure) {
+ return ProcedureType.ALTER_LOGICAL_VIEW_PROCEDURE;
}
return null;
}
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java
index 7e6ffa008ea..31bb30aa7c4 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/procedure/store/ProcedureType.java
@@ -78,7 +78,8 @@ public enum ProcedureType {
PIPE_HANDLE_META_CHANGE_PROCEDURE((short) 1102),
/** logical view */
- DELETE_LOGICAL_VIEW_PROCEDURE((short) 1200);
+ DELETE_LOGICAL_VIEW_PROCEDURE((short) 1200),
+ ALTER_LOGICAL_VIEW_PROCEDURE((short) 12001);
private final short typeCode;
diff --git
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
index 6e73fee6675..7b657c1080d 100644
---
a/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
+++
b/confignode/src/main/java/org/apache/iotdb/confignode/service/thrift/ConfigNodeRPCServiceProcessor.java
@@ -66,6 +66,7 @@ import org.apache.iotdb.confignode.manager.ConfigManager;
import org.apache.iotdb.confignode.manager.consensus.ConsensusManager;
import org.apache.iotdb.confignode.rpc.thrift.IConfigNodeRPCService;
import org.apache.iotdb.confignode.rpc.thrift.TAddConsensusGroupReq;
+import org.apache.iotdb.confignode.rpc.thrift.TAlterLogicalViewReq;
import org.apache.iotdb.confignode.rpc.thrift.TAlterSchemaTemplateReq;
import org.apache.iotdb.confignode.rpc.thrift.TAuthorizerReq;
import org.apache.iotdb.confignode.rpc.thrift.TAuthorizerResp;
@@ -863,6 +864,11 @@ public class ConfigNodeRPCServiceProcessor implements
IConfigNodeRPCService.Ifac
return configManager.deleteLogicalView(req);
}
+ @Override
+ public TSStatus alterLogicalView(TAlterLogicalViewReq req) throws TException
{
+ return configManager.alterLogicalView(req);
+ }
+
@Override
@Deprecated
public TSStatus createPipeSink(TPipeSinkInfo req) {
diff --git
a/integration-test/src/test/java/org/apache/iotdb/db/it/schema/view/IoTDBAliasSeriesIT.java
b/integration-test/src/test/java/org/apache/iotdb/db/it/schema/view/IoTDBAliasSeriesIT.java
index fa66c77ce80..2045f23b59e 100644
---
a/integration-test/src/test/java/org/apache/iotdb/db/it/schema/view/IoTDBAliasSeriesIT.java
+++
b/integration-test/src/test/java/org/apache/iotdb/db/it/schema/view/IoTDBAliasSeriesIT.java
@@ -24,8 +24,9 @@ import org.apache.iotdb.itbase.category.ClusterIT;
import org.apache.iotdb.itbase.category.LocalStandaloneIT;
import org.junit.After;
+import org.junit.AfterClass;
import org.junit.Assert;
-import org.junit.Before;
+import org.junit.BeforeClass;
import org.junit.Test;
import org.junit.experimental.categories.Category;
import org.junit.runner.RunWith;
@@ -39,11 +40,16 @@ import java.sql.Statement;
@Category({LocalStandaloneIT.class, ClusterIT.class})
public class IoTDBAliasSeriesIT {
- @Before
- public void setUp() throws Exception {
+ @BeforeClass
+ public static void setUpCluster() throws Exception {
EnvFactory.getEnv().initClusterEnvironment();
}
+ @AfterClass
+ public static void tearDownCluster() throws Exception {
+ EnvFactory.getEnv().cleanClusterEnvironment();
+ }
+
@After
public void tearDown() throws Exception {
try (Connection connection = EnvFactory.getEnv().getConnection();
@@ -54,7 +60,6 @@ public class IoTDBAliasSeriesIT {
// If database is null, it will throw exception. Do nothing.
}
}
- EnvFactory.getEnv().cleanClusterEnvironment();
}
@Test
diff --git
a/integration-test/src/test/java/org/apache/iotdb/db/it/schema/view/IoTDBAlterViewIT.java
b/integration-test/src/test/java/org/apache/iotdb/db/it/schema/view/IoTDBAlterViewIT.java
new file mode 100644
index 00000000000..698535c82b2
--- /dev/null
+++
b/integration-test/src/test/java/org/apache/iotdb/db/it/schema/view/IoTDBAlterViewIT.java
@@ -0,0 +1,117 @@
+/*
+ * 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.iotdb.db.it.schema.view;
+
+import org.apache.iotdb.it.env.EnvFactory;
+import org.apache.iotdb.it.framework.IoTDBTestRunner;
+import org.apache.iotdb.itbase.category.ClusterIT;
+import org.apache.iotdb.itbase.category.LocalStandaloneIT;
+
+import org.junit.After;
+import org.junit.AfterClass;
+import org.junit.Assert;
+import org.junit.BeforeClass;
+import org.junit.Test;
+import org.junit.experimental.categories.Category;
+import org.junit.runner.RunWith;
+
+import java.sql.Connection;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.sql.Statement;
+
+@RunWith(IoTDBTestRunner.class)
+@Category({LocalStandaloneIT.class, ClusterIT.class})
+public class IoTDBAlterViewIT {
+
+ @BeforeClass
+ public static void setUpCluster() throws Exception {
+ EnvFactory.getEnv().initClusterEnvironment();
+ }
+
+ @AfterClass
+ public static void tearDownCluster() throws Exception {
+ EnvFactory.getEnv().cleanClusterEnvironment();
+ }
+
+ @After
+ public void tearDown() throws Exception {
+ try (Connection connection = EnvFactory.getEnv().getConnection();
+ Statement statement = connection.createStatement()) {
+ try {
+ statement.execute("DELETE DATABASE root.**");
+ } catch (Exception e) {
+ // If database is null, it will throw exception. Do nothing.
+ }
+ }
+ }
+
+ @Test
+ public void testAlterView() throws SQLException {
+ try (Connection connection = EnvFactory.getEnv().getConnection();
+ Statement statement = connection.createStatement()) {
+ statement.execute("create timeseries root.db.d1.s1 with datatype=INT32");
+ statement.execute("create timeseries root.db.d1.s2 with datatype=INT32");
+ statement.execute("create timeseries root.db.d1.s3 with datatype=INT32");
+ statement.execute("create timeseries root.db.d2.s1 with datatype=INT32");
+ statement.execute("create timeseries root.db.d2.s2 with datatype=INT32");
+ statement.execute("create timeseries root.db.d2.s3 with datatype=INT32");
+
+ statement.execute(
+ "create view root(view.d1.s1, view.d1.s2, view.d1.s3, view.d2.s1,
view.d2.s2, view.d2.s3) as root(db.d1.s1, db.d1.s2, db.d1.s3, db.d2.s1,
db.d2.s2, db.d2.s3)");
+
+ String[][] map =
+ new String[][] {
+ new String[] {"root.view.d1.s1", "root.db.d1.s1"},
+ new String[] {"root.view.d1.s2", "root.db.d1.s2"},
+ new String[] {"root.view.d1.s3", "root.db.d1.s3"},
+ new String[] {"root.view.d2.s1", "root.db.d2.s1"},
+ new String[] {"root.view.d2.s2", "root.db.d2.s2"},
+ new String[] {"root.view.d2.s3", "root.db.d2.s3"},
+ };
+ for (String[] strings : map) {
+ try (ResultSet resultSet =
+ statement.executeQuery(String.format("show view %s", strings[0])))
{
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(strings[1], resultSet.getString("Source"));
+ }
+ }
+
+ statement.execute(
+ "alter view root(view.d1.s1, view.d1.s2, view.d1.s3, view.d2.s1,
view.d2.s2, view.d2.s3) as root(db.d2.s2, db.d2.s3, db.d2.s1, db.d1.s2,
db.d1.s3, db.d1.s1)");
+
+ map =
+ new String[][] {
+ new String[] {"root.view.d1.s1", "root.db.d2.s2"},
+ new String[] {"root.view.d1.s2", "root.db.d2.s3"},
+ new String[] {"root.view.d1.s3", "root.db.d2.s1"},
+ new String[] {"root.view.d2.s1", "root.db.d1.s2"},
+ new String[] {"root.view.d2.s2", "root.db.d1.s3"},
+ new String[] {"root.view.d2.s3", "root.db.d1.s1"},
+ };
+ for (String[] strings : map) {
+ try (ResultSet resultSet =
+ statement.executeQuery(String.format("show view %s", strings[0])))
{
+ Assert.assertTrue(resultSet.next());
+ Assert.assertEquals(strings[1], resultSet.getString("Source"));
+ }
+ }
+ }
+ }
+}
diff --git a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
index e9d41faea32..66b402ba5b7 100644
--- a/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
+++ b/iotdb-protocol/thrift-confignode/src/main/thrift/confignode.thrift
@@ -664,6 +664,11 @@ struct TDeleteLogicalViewReq{
2: required binary pathPatternTree
}
+struct TAlterLogicalViewReq{
+ 1: required string queryId
+ 2: required binary viewBinary
+}
+
// ====================================================
// CQ
// ====================================================
@@ -1267,6 +1272,8 @@ service IConfigNodeRPCService {
common.TSStatus deleteLogicalView(TDeleteLogicalViewReq req)
+ common.TSStatus alterLogicalView(TAlterLogicalViewReq req)
+
// ======================================================
// Sync
// ======================================================
diff --git a/iotdb-protocol/thrift/src/main/thrift/datanode.thrift
b/iotdb-protocol/thrift/src/main/thrift/datanode.thrift
index 727d2e8f2a1..43d26a03835 100644
--- a/iotdb-protocol/thrift/src/main/thrift/datanode.thrift
+++ b/iotdb-protocol/thrift/src/main/thrift/datanode.thrift
@@ -402,6 +402,11 @@ struct TDeleteViewSchemaReq{
2: required binary pathPatternTree
}
+struct TAlterViewReq{
+ 1: required list<common.TConsensusGroupId> schemaRegionIdList
+ 2: required list<binary> viewBinaryList
+}
+
// ====================================================
// CQ
// ====================================================
@@ -785,6 +790,8 @@ service IDataNodeRPCService {
common.TSStatus deleteViewSchema(TDeleteViewSchemaReq req)
+ common.TSStatus alterView(TAlterViewReq req)
+
/**
* Send pipeMetas to DataNodes, for synchronization
*/
diff --git
a/server/src/main/java/org/apache/iotdb/db/client/ConfigNodeClient.java
b/server/src/main/java/org/apache/iotdb/db/client/ConfigNodeClient.java
index f9d69c2caa4..1e57b845366 100644
--- a/server/src/main/java/org/apache/iotdb/db/client/ConfigNodeClient.java
+++ b/server/src/main/java/org/apache/iotdb/db/client/ConfigNodeClient.java
@@ -35,6 +35,7 @@ import
org.apache.iotdb.commons.client.sync.SyncThriftClientWithErrorHandler;
import org.apache.iotdb.commons.consensus.ConfigRegionId;
import org.apache.iotdb.confignode.rpc.thrift.IConfigNodeRPCService;
import org.apache.iotdb.confignode.rpc.thrift.TAddConsensusGroupReq;
+import org.apache.iotdb.confignode.rpc.thrift.TAlterLogicalViewReq;
import org.apache.iotdb.confignode.rpc.thrift.TAlterSchemaTemplateReq;
import org.apache.iotdb.confignode.rpc.thrift.TAuthorizerReq;
import org.apache.iotdb.confignode.rpc.thrift.TAuthorizerResp;
@@ -1717,6 +1718,27 @@ public class ConfigNodeClient implements
IConfigNodeRPCService.Iface, ThriftClie
throw new TException(MSG_RECONNECTION_FAIL);
}
+ @Override
+ public TSStatus alterLogicalView(TAlterLogicalViewReq req) throws TException
{
+ for (int i = 0; i < RETRY_NUM; i++) {
+ try {
+ TSStatus status = client.alterLogicalView(req);
+ if (!updateConfigNodeLeader(status)) {
+ return status;
+ }
+ } catch (TException e) {
+ logger.warn(
+ "Failed to connect to ConfigNode {} from DataNode {} when
executing {}",
+ configNode,
+ config.getAddressAndPort(),
+ Thread.currentThread().getStackTrace()[1].getMethodName());
+ configLeader = null;
+ }
+ waitAndReconnect();
+ }
+ throw new TException(MSG_RECONNECTION_FAIL);
+ }
+
@Override
public TSStatus createPipeSink(TPipeSinkInfo req) throws TException {
for (int i = 0; i < RETRY_NUM; i++) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/exception/metadata/view/ViewNotExistException.java
b/server/src/main/java/org/apache/iotdb/db/exception/metadata/view/ViewNotExistException.java
index 4bcd29c4071..e991640da64 100644
---
a/server/src/main/java/org/apache/iotdb/db/exception/metadata/view/ViewNotExistException.java
+++
b/server/src/main/java/org/apache/iotdb/db/exception/metadata/view/ViewNotExistException.java
@@ -27,25 +27,12 @@ public class ViewNotExistException extends
MetadataException {
private static final String VIEW_NOT_EXIST_WRONG_MESSAGE = "View [%s] does
not exist";
- private static final String NORMAL_VIEW_NOT_EXIST_WRONG_MESSAGE =
- "View [%s] does not exist or is represented by schema template";
-
- private static final String TEMPLATE_VIEW_NOT_EXIST_WRONG_MESSAGE =
- "View [%s] does not exist or is not represented by schema template";
-
public ViewNotExistException(String path) {
super(
String.format(VIEW_NOT_EXIST_WRONG_MESSAGE, path),
TSStatusCode.PATH_NOT_EXIST.getStatusCode());
}
- public ViewNotExistException(String path, boolean isUserException) {
- super(
- String.format(VIEW_NOT_EXIST_WRONG_MESSAGE, path),
- TSStatusCode.PATH_NOT_EXIST.getStatusCode(),
- isUserException);
- }
-
public ViewNotExistException(List<String> paths) {
super(
String.format(
@@ -55,16 +42,4 @@ public class ViewNotExistException extends MetadataException
{
: paths.get(0) + " ... " + paths.get(paths.size() - 1)),
TSStatusCode.PATH_NOT_EXIST.getStatusCode());
}
-
- public ViewNotExistException(List<String> paths, boolean isTemplateSeries) {
- super(
- String.format(
- isTemplateSeries
- ? TEMPLATE_VIEW_NOT_EXIST_WRONG_MESSAGE
- : NORMAL_VIEW_NOT_EXIST_WRONG_MESSAGE,
- paths.size() == 1
- ? paths.get(0)
- : paths.get(0) + " ... " + paths.get(paths.size() - 1)),
- TSStatusCode.PATH_NOT_EXIST.getStatusCode());
- }
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanType.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanType.java
index 33e4dbd2136..d8a89365b4e 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanType.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanType.java
@@ -49,6 +49,7 @@ public enum SchemaRegionPlanType {
PRE_DELETE_LOGICAL_VIEW((byte) 67),
ROLLBACK_PRE_DELETE_LOGICAL_VIEW((byte) 68),
DELETE_LOGICAL_VIEW((byte) 69),
+ ALTER_LOGICAL_VIEW((byte) 70),
// query plan doesn't need any ser/deSer, thus use one type to represent all
READ_SCHEMA(Byte.MAX_VALUE);
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanVisitor.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanVisitor.java
index 4936e3160aa..e401c171b65 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanVisitor.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/SchemaRegionPlanVisitor.java
@@ -32,6 +32,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTempla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeleteTimeSeriesPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeleteTimeSeriesPlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IPreDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IRollbackPreDeleteLogicalViewPlan;
@@ -98,6 +99,10 @@ public abstract class SchemaRegionPlanVisitor<R, C> {
return visitSchemaRegionPlan(createLogicalViewPlan, context);
}
+ public R visitAlterLogicalView(IAlterLogicalViewPlan alterLogicalViewPlan, C
context) {
+ return visitSchemaRegionPlan(alterLogicalViewPlan, context);
+ }
+
public R visitPreDeleteLogicalView(
IPreDeleteLogicalViewPlan preDeleteLogicalViewPlan, C context) {
return visitSchemaRegionPlan(preDeleteLogicalViewPlan, context);
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanDeserializer.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanDeserializer.java
index e4ff058041e..a8b68c97a27 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanDeserializer.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanDeserializer.java
@@ -41,6 +41,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTempla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeleteTimeSeriesPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeleteTimeSeriesPlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IPreDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IRollbackPreDeleteLogicalViewPlan;
@@ -386,5 +387,13 @@ public class SchemaRegionPlanDeserializer implements
IDeserializer<ISchemaRegion
deleteLogicalViewPlan.setPath((PartialPath)
PathDeserializeUtil.deserialize(buffer));
return deleteLogicalViewPlan;
}
+
+ @Override
+ public ISchemaRegionPlan visitAlterLogicalView(
+ IAlterLogicalViewPlan alterLogicalViewPlan, ByteBuffer buffer) {
+ alterLogicalViewPlan.setViewPath((PartialPath)
PathDeserializeUtil.deserialize(buffer));
+
alterLogicalViewPlan.setSourceExpression(ViewExpression.deserialize(buffer));
+ return alterLogicalViewPlan;
+ }
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanSerializer.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanSerializer.java
index d4dd9041b40..bdaaa6e1f69 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanSerializer.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanSerializer.java
@@ -37,6 +37,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTempla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeleteTimeSeriesPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeleteTimeSeriesPlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IPreDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IRollbackPreDeleteLogicalViewPlan;
@@ -461,5 +462,17 @@ public class SchemaRegionPlanSerializer implements
ISerializer<ISchemaRegionPlan
return new SchemaRegionPlanSerializationResult(e);
}
}
+
+ @Override
+ public SchemaRegionPlanSerializationResult visitAlterLogicalView(
+ IAlterLogicalViewPlan alterLogicalViewPlan, DataOutputStream stream) {
+ try {
+ alterLogicalViewPlan.getViewPath().serialize(stream);
+ ViewExpression.serialize(alterLogicalViewPlan.getSourceExpression(),
stream);
+ return SchemaRegionPlanSerializationResult.SUCCESS;
+ } catch (IOException e) {
+ return new SchemaRegionPlanSerializationResult(e);
+ }
+ }
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanTxtSerializer.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanTxtSerializer.java
index a4445f967de..9ad3296c998 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanTxtSerializer.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/SchemaRegionPlanTxtSerializer.java
@@ -37,6 +37,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTempla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeleteTimeSeriesPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeleteTimeSeriesPlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IPreDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IRollbackPreDeleteLogicalViewPlan;
@@ -286,5 +287,15 @@ public class SchemaRegionPlanTxtSerializer implements
ISerializer<ISchemaRegionP
stringBuilder.append(deleteLogicalViewPlan.getPath().getFullPath());
return null;
}
+
+ @Override
+ public Void visitAlterLogicalView(
+ IAlterLogicalViewPlan alterLogicalViewPlan, StringBuilder
stringBuilder) {
+ stringBuilder
+ .append(alterLogicalViewPlan.getViewPath())
+ .append(FIELD_SEPARATOR)
+ .append(alterLogicalViewPlan.getSourceExpression().toString());
+ return null;
+ }
}
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/write/AlterLogicalViewPlanImpl.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/write/AlterLogicalViewPlanImpl.java
new file mode 100644
index 00000000000..78bf57ef16d
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/write/AlterLogicalViewPlanImpl.java
@@ -0,0 +1,56 @@
+/*
+ * 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.iotdb.db.metadata.plan.schemaregion.impl.write;
+
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
+
+public class AlterLogicalViewPlanImpl implements IAlterLogicalViewPlan {
+
+ private PartialPath targetPath = null;
+ private ViewExpression sourceExpression = null;
+
+ public AlterLogicalViewPlanImpl() {}
+
+ public AlterLogicalViewPlanImpl(PartialPath targetPath, ViewExpression
sourceExpression) {
+ this.targetPath = targetPath;
+ this.sourceExpression = sourceExpression;
+ }
+
+ @Override
+ public PartialPath getViewPath() {
+ return targetPath;
+ }
+
+ public ViewExpression getSourceExpression() {
+ return this.sourceExpression;
+ }
+
+ @Override
+ public void setViewPath(PartialPath path) {
+ this.targetPath = path;
+ }
+
+ @Override
+ public void setSourceExpression(ViewExpression viewExpression) {
+ this.sourceExpression = viewExpression;
+ }
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/write/SchemaRegionWritePlanFactory.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/write/SchemaRegionWritePlanFactory.java
index 1ee9a2bc834..1ed86908eb6 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/write/SchemaRegionWritePlanFactory.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/impl/write/SchemaRegionWritePlanFactory.java
@@ -36,6 +36,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTempla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeleteTimeSeriesPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeleteTimeSeriesPlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IPreDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IRollbackPreDeleteLogicalViewPlan;
@@ -84,6 +85,8 @@ public class SchemaRegionWritePlanFactory {
return new RollbackPreDeleteLogicalViewPlanImpl();
case DELETE_LOGICAL_VIEW:
return new DeleteLogicalViewPlanImpl();
+ case ALTER_LOGICAL_VIEW:
+ return new AlterLogicalViewPlanImpl();
default:
throw new UnsupportedOperationException(
String.format(
@@ -187,4 +190,9 @@ public class SchemaRegionWritePlanFactory {
public static IDeleteLogicalViewPlan getDeleteLogicalViewPlan(PartialPath
path) {
return new DeleteLogicalViewPlanImpl(path);
}
+
+ public static IAlterLogicalViewPlan getAlterLogicalViewPlan(
+ PartialPath targetPath, ViewExpression sourceExpression) {
+ return new AlterLogicalViewPlanImpl(targetPath, sourceExpression);
+ }
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/write/view/IAlterLogicalViewPlan.java
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/write/view/IAlterLogicalViewPlan.java
new file mode 100644
index 00000000000..6f06ecb0855
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/plan/schemaregion/write/view/IAlterLogicalViewPlan.java
@@ -0,0 +1,46 @@
+/*
+ * 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.iotdb.db.metadata.plan.schemaregion.write.view;
+
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
+import org.apache.iotdb.db.metadata.plan.schemaregion.ISchemaRegionPlan;
+import org.apache.iotdb.db.metadata.plan.schemaregion.SchemaRegionPlanType;
+import org.apache.iotdb.db.metadata.plan.schemaregion.SchemaRegionPlanVisitor;
+
+public interface IAlterLogicalViewPlan extends ISchemaRegionPlan {
+ @Override
+ default SchemaRegionPlanType getPlanType() {
+ return SchemaRegionPlanType.ALTER_LOGICAL_VIEW;
+ }
+
+ @Override
+ default <R, C> R accept(SchemaRegionPlanVisitor<R, C> visitor, C context) {
+ return visitor.visitAlterLogicalView(this, context);
+ }
+
+ PartialPath getViewPath();
+
+ ViewExpression getSourceExpression();
+
+ void setViewPath(PartialPath path);
+
+ void setSourceExpression(ViewExpression viewExpression);
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/ISchemaRegion.java
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/ISchemaRegion.java
index fb204799738..3c368e7d124 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/ISchemaRegion.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/ISchemaRegion.java
@@ -37,6 +37,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.ICreateTimeSeriesPla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import org.apache.iotdb.db.metadata.query.info.IDeviceSchemaInfo;
import org.apache.iotdb.db.metadata.query.info.INodeSchemaInfo;
import org.apache.iotdb.db.metadata.query.info.ITimeSeriesSchemaInfo;
@@ -183,6 +184,8 @@ public interface ISchemaRegion {
void deleteLogicalView(PathPatternTree patternTree) throws MetadataException;
+ void alterLogicalView(IAlterLogicalViewPlan alterLogicalViewPlan) throws
MetadataException;
+
// endregion
// region Interfaces for metadata info Query
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
index d4ac538c03f..26bb62afd68 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionMemoryImpl.java
@@ -29,6 +29,7 @@ import
org.apache.iotdb.commons.schema.ClusterSchemaQuotaLevel;
import org.apache.iotdb.commons.schema.filter.SchemaFilterType;
import org.apache.iotdb.commons.schema.node.role.IDeviceMNode;
import org.apache.iotdb.commons.schema.node.role.IMeasurementMNode;
+import org.apache.iotdb.commons.schema.view.LogicalViewSchema;
import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
import org.apache.iotdb.commons.utils.FileUtils;
import org.apache.iotdb.consensus.ConsensusFactory;
@@ -69,6 +70,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTempla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeleteTimeSeriesPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeleteTimeSeriesPlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IPreDeleteLogicalViewPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IRollbackPreDeleteLogicalViewPlan;
@@ -839,6 +841,27 @@ public class SchemaRegionMemoryImpl implements
ISchemaRegion {
}
}
+ @Override
+ public void alterLogicalView(IAlterLogicalViewPlan alterLogicalViewPlan)
+ throws MetadataException {
+ IMeasurementMNode<IMemMNode> leafMNode =
+ mtree.getMeasurementMNode(alterLogicalViewPlan.getViewPath());
+ if (!leafMNode.isLogicalView()) {
+ throw new MetadataException(
+ String.format("[%s] is no view.",
alterLogicalViewPlan.getViewPath()));
+ }
+ leafMNode.setSchema(
+ new LogicalViewSchema(leafMNode.getName(),
alterLogicalViewPlan.getSourceExpression()));
+ // write log
+ if (!isRecovering) {
+ try {
+ writeToMLog(alterLogicalViewPlan);
+ } catch (IOException e) {
+ throw new MetadataException(e);
+ }
+ }
+ }
+
private void deleteSingleTimeseriesInBlackList(PartialPath path)
throws MetadataException, IOException {
IMeasurementMNode<IMemMNode> measurementMNode =
mtree.deleteTimeseries(path);
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
index 041f88a8a72..95ec2f4e10c 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/schemaregion/SchemaRegionSchemaFileImpl.java
@@ -71,6 +71,7 @@ import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeactivateTempla
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IPreDeleteTimeSeriesPlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeactivateTemplatePlan;
import
org.apache.iotdb.db.metadata.plan.schemaregion.write.IRollbackPreDeleteTimeSeriesPlan;
+import
org.apache.iotdb.db.metadata.plan.schemaregion.write.view.IAlterLogicalViewPlan;
import org.apache.iotdb.db.metadata.query.info.IDeviceSchemaInfo;
import org.apache.iotdb.db.metadata.query.info.INodeSchemaInfo;
import org.apache.iotdb.db.metadata.query.info.ITimeSeriesSchemaInfo;
@@ -891,6 +892,12 @@ public class SchemaRegionSchemaFileImpl implements
ISchemaRegion {
throw new UnsupportedOperationException();
}
+ @Override
+ public void alterLogicalView(IAlterLogicalViewPlan alterLogicalViewPlan)
+ throws MetadataException {
+ throw new UnsupportedOperationException();
+ }
+
private void deleteSingleTimeseriesInBlackList(PartialPath path)
throws MetadataException, IOException {
IMeasurementMNode<ICachedMNode> measurementMNode =
mtree.deleteTimeseries(path);
diff --git
a/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
b/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
index 8b9db82b733..515001e1af9 100644
---
a/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
+++
b/server/src/main/java/org/apache/iotdb/db/metadata/visitor/SchemaExecutionVisitor.java
@@ -51,6 +51,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.Measurement
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.PreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackPreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackSchemaBlackListNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.AlterLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.ConstructLogicalViewBlackListNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.CreateLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.DeleteLogicalViewNode;
@@ -465,6 +466,25 @@ public class SchemaExecutionVisitor extends
PlanVisitor<TSStatus, ISchemaRegion>
return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS, "Execute
successfully");
}
+ @Override
+ public TSStatus visitAlterLogicalView(AlterLogicalViewNode node,
ISchemaRegion schemaRegion) {
+ Map<PartialPath, ViewExpression> viewPathToSourceMap =
node.getViewPathToSourceMap();
+ List<TSStatus> failingStatus = new ArrayList<>();
+ for (Map.Entry<PartialPath, ViewExpression> entry :
viewPathToSourceMap.entrySet()) {
+ try {
+ schemaRegion.alterLogicalView(
+
SchemaRegionWritePlanFactory.getAlterLogicalViewPlan(entry.getKey(),
entry.getValue()));
+ } catch (MetadataException e) {
+ logger.error("{}: MetaData error: ", IoTDBConstant.GLOBAL_DB_NAME, e);
+ failingStatus.add(RpcUtils.getStatus(e.getErrorCode(),
e.getMessage()));
+ }
+ }
+ if (!failingStatus.isEmpty()) {
+ return RpcUtils.getStatus(failingStatus);
+ }
+ return RpcUtils.getStatus(TSStatusCode.SUCCESS_STATUS, "Execute
successfully");
+ }
+
@Override
public TSStatus visitConstructLogicalViewBlackList(
ConstructLogicalViewBlackListNode node, ISchemaRegion schemaRegion) {
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/config/executor/ClusterConfigTaskExecutor.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/config/executor/ClusterConfigTaskExecutor.java
index 53061ae18df..1136c31c55e 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/config/executor/ClusterConfigTaskExecutor.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/execution/config/executor/ClusterConfigTaskExecutor.java
@@ -45,6 +45,7 @@ import
org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
import org.apache.iotdb.commons.trigger.service.TriggerExecutableManager;
import org.apache.iotdb.commons.udf.service.UDFClassLoader;
import org.apache.iotdb.commons.udf.service.UDFExecutableManager;
+import org.apache.iotdb.confignode.rpc.thrift.TAlterLogicalViewReq;
import org.apache.iotdb.confignode.rpc.thrift.TAlterSchemaTemplateReq;
import org.apache.iotdb.confignode.rpc.thrift.TCountDatabaseResp;
import org.apache.iotdb.confignode.rpc.thrift.TCountTimeSlotListReq;
@@ -148,6 +149,7 @@ import
org.apache.iotdb.db.mpp.plan.execution.config.sys.quota.ShowSpaceQuotaTas
import
org.apache.iotdb.db.mpp.plan.execution.config.sys.quota.ShowThrottleQuotaTask;
import org.apache.iotdb.db.mpp.plan.execution.config.sys.sync.ShowPipeSinkTask;
import org.apache.iotdb.db.mpp.plan.expression.Expression;
+import
org.apache.iotdb.db.mpp.plan.expression.visitor.TransformToViewExpressionVisitor;
import org.apache.iotdb.db.mpp.plan.statement.metadata.CountDatabaseStatement;
import
org.apache.iotdb.db.mpp.plan.statement.metadata.CountTimeSlotListStatement;
import
org.apache.iotdb.db.mpp.plan.statement.metadata.CreateContinuousQueryStatement;
@@ -201,6 +203,7 @@ import org.apache.iotdb.rpc.StatementExecutionException;
import org.apache.iotdb.rpc.TSStatusCode;
import org.apache.iotdb.trigger.api.Trigger;
import org.apache.iotdb.trigger.api.enums.FailureStrategy;
+import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
import org.apache.iotdb.udf.api.UDTF;
import com.google.common.util.concurrent.SettableFuture;
@@ -1857,17 +1860,43 @@ public class ClusterConfigTaskExecutor implements
IConfigTaskExecutor {
public SettableFuture<ConfigTaskResult> alterLogicalView(
String queryId, AlterLogicalViewStatement alterLogicalViewStatement) {
SettableFuture<ConfigTaskResult> future = SettableFuture.create();
- // delete old view
- TDeleteLogicalViewReq req =
- new TDeleteLogicalViewReq(
- queryId,
-
serializePatternListToByteBuffer(alterLogicalViewStatement.getTargetPathList()));
+ CreateLogicalViewStatement createLogicalViewStatement = new
CreateLogicalViewStatement();
+
createLogicalViewStatement.setTargetPaths(alterLogicalViewStatement.getTargetPaths());
+
createLogicalViewStatement.setSourcePaths(alterLogicalViewStatement.getSourcePaths());
+
createLogicalViewStatement.setQueryStatement(alterLogicalViewStatement.getQueryStatement());
+
+ Analyzer.validate(createLogicalViewStatement);
+
+ // Transform all Expressions into ViewExpressions.
+ TransformToViewExpressionVisitor transformToViewExpressionVisitor =
+ new TransformToViewExpressionVisitor();
+ List<Expression> expressionList =
createLogicalViewStatement.getSourceExpressionList();
+ List<ViewExpression> viewExpressionList = new ArrayList<>();
+ for (Expression expression : expressionList) {
+
viewExpressionList.add(transformToViewExpressionVisitor.process(expression,
null));
+ }
+
+ List<PartialPath> viewPathList =
createLogicalViewStatement.getTargetPathList();
+
+ ByteArrayOutputStream stream = new ByteArrayOutputStream();
+ try {
+ ReadWriteIOUtils.write(viewPathList.size(), stream);
+ for (int i = 0; i < viewPathList.size(); i++) {
+ viewPathList.get(i).serialize(stream);
+ ViewExpression.serialize(viewExpressionList.get(i), stream);
+ }
+ } catch (IOException e) {
+ throw new RuntimeException(e);
+ }
+
+ TAlterLogicalViewReq req =
+ new TAlterLogicalViewReq(queryId,
ByteBuffer.wrap(stream.toByteArray()));
try (ConfigNodeClient client =
CLUSTER_DELETION_CONFIG_NODE_CLIENT_MANAGER.borrowClient(ConfigNodeInfo.CONFIG_REGION_ID))
{
TSStatus tsStatus;
do {
try {
- tsStatus = client.deleteLogicalView(req);
+ tsStatus = client.alterLogicalView(req);
} catch (TTransportException e) {
if (e.getType() == TTransportException.TIMED_OUT
|| e.getCause() instanceof SocketTimeoutException) {
@@ -1882,43 +1911,18 @@ public class ClusterConfigTaskExecutor implements
IConfigTaskExecutor {
if (TSStatusCode.SUCCESS_STATUS.getStatusCode() != tsStatus.getCode()) {
LOGGER.warn(
- "Failed to execute delete view {}, status is {}.",
+ "Failed to execute alter view {}, status is {}.",
alterLogicalViewStatement.getTargetPathList(),
tsStatus);
future.setException(new IoTDBException(tsStatus.getMessage(),
tsStatus.getCode()));
- return future;
+ } else {
+ future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS));
}
+ return future;
} catch (ClientManagerException | TException e) {
future.setException(e);
return future;
}
-
- // recreate the logical view
- CreateLogicalViewStatement createLogicalViewStatement = new
CreateLogicalViewStatement();
-
createLogicalViewStatement.setTargetPaths(alterLogicalViewStatement.getTargetPaths());
-
createLogicalViewStatement.setSourcePaths(alterLogicalViewStatement.getSourcePaths());
- createLogicalViewStatement.setSourceQueryStatement(
- alterLogicalViewStatement.getQueryStatement());
-
- ExecutionResult executionResult =
- Coordinator.getInstance()
- .execute(
- createLogicalViewStatement,
- 0,
- null,
- "",
- ClusterPartitionFetcher.getInstance(),
- ClusterSchemaFetcher.getInstance(),
-
IoTDBDescriptor.getInstance().getConfig().getQueryTimeoutThreshold());
- if (executionResult.status.getCode() !=
TSStatusCode.SUCCESS_STATUS.getStatusCode()) {
- future.setException(
- new IoTDBException(
- executionResult.status.getMessage(),
executionResult.status.getCode()));
- } else {
- future.set(new ConfigTaskResult(TSStatusCode.SUCCESS_STATUS));
- }
-
- return future;
}
@Override
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/ASTVisitor.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/ASTVisitor.java
index 69c0b704059..34cd5711ea8 100644
--- a/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/ASTVisitor.java
+++ b/server/src/main/java/org/apache/iotdb/db/mpp/plan/parser/ASTVisitor.java
@@ -1072,10 +1072,12 @@ public class ASTVisitor extends
IoTDBSqlParserBaseVisitor<Statement> {
ctx.viewTargetPaths(),
alterLogicalViewStatement::setTargetFullPaths,
alterLogicalViewStatement::setTargetPathsGroup,
- alterLogicalViewStatement::setTargetIntoItem);
- if (alterLogicalViewStatement.getIntoItem() != null) {
- throw new SemanticException("Can not use char '$' or into item in
alter view statement.");
- }
+ intoItem -> {
+ if (intoItem != null) {
+ throw new SemanticException(
+ "Can not use char '$' or into item in alter view
statement.");
+ }
+ });
// parse source
parseViewSourcePaths(
ctx.viewSourcePaths(),
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
index ecda02923eb..21328034a23 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/LogicalPlanVisitor.java
@@ -834,11 +834,11 @@ public class LogicalPlanVisitor extends
StatementVisitor<PlanNode, MPPQueryConte
@Override
public PlanNode visitCreateLogicalView(
CreateLogicalViewStatement createLogicalViewStatement, MPPQueryContext
context) {
- // Transform all Expressions into ViewExpressions.
- TransformToViewExpressionVisitor transformToViewExpressionVisitor =
- new TransformToViewExpressionVisitor();
List<ViewExpression> viewExpressionList = new ArrayList<>();
if (createLogicalViewStatement.getViewExpression() == null) {
+ // Transform all Expressions into ViewExpressions.
+ TransformToViewExpressionVisitor transformToViewExpressionVisitor =
+ new TransformToViewExpressionVisitor();
List<Expression> expressionList =
createLogicalViewStatement.getSourceExpressionList();
for (Expression expression : expressionList) {
viewExpressionList.add(transformToViewExpressionVisitor.process(expression,
null));
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
index f0ff085eded..a6faf620f5f 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanNodeType.java
@@ -51,6 +51,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.InvalidateS
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.PreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackPreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackSchemaBlackListNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.AlterLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.ConstructLogicalViewBlackListNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.CreateLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.DeleteLogicalViewNode;
@@ -177,7 +178,8 @@ public enum PlanNodeType {
CONSTRUCT_LOGICAL_VIEW_BLACK_LIST((short) 74),
ROLLBACK_LOGICAL_VIEW_BLACK_LIST((short) 75),
DELETE_LOGICAL_VIEW((short) 76),
- LOGICAL_VIEW_SCHEMA_SCAN((short) 77);
+ LOGICAL_VIEW_SCHEMA_SCAN((short) 77),
+ ALTER_LOGICAL_VIEW((short) 78);
public static final int BYTES = Short.BYTES;
@@ -380,6 +382,8 @@ public enum PlanNodeType {
return DeleteLogicalViewNode.deserialize(buffer);
case 77:
return LogicalViewSchemaScanNode.deserialize(buffer);
+ case 78:
+ return AlterLogicalViewNode.deserialize(buffer);
default:
throw new IllegalArgumentException("Invalid node type: " + nodeType);
}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
index d2955839a8c..f30f11d427d 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/PlanVisitor.java
@@ -48,6 +48,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.InternalCre
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.PreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackPreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackSchemaBlackListNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.AlterLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.ConstructLogicalViewBlackListNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.CreateLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.DeleteLogicalViewNode;
@@ -394,6 +395,10 @@ public abstract class PlanVisitor<R, C> {
return visitPlan(node, context);
}
+ public R visitAlterLogicalView(AlterLogicalViewNode node, C context) {
+ return visitPlan(node, context);
+ }
+
/////////////////////////////////////////////////////////////////////////////////////////////////
// Data Write Node
/////////////////////////////////////////////////////////////////////////////////////////////////
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/write/view/AlterLogicalViewNode.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/write/view/AlterLogicalViewNode.java
new file mode 100644
index 00000000000..5cdde4a5378
--- /dev/null
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/planner/plan/node/metedata/write/view/AlterLogicalViewNode.java
@@ -0,0 +1,186 @@
+/*
+ * 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.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view;
+
+import org.apache.iotdb.common.rpc.thrift.TRegionReplicaSet;
+import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.path.PathDeserializeUtil;
+import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
+import org.apache.iotdb.db.mpp.plan.analyze.Analysis;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNode;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeId;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanNodeType;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.PlanVisitor;
+import org.apache.iotdb.db.mpp.plan.planner.plan.node.WritePlanNode;
+import org.apache.iotdb.tsfile.exception.NotImplementedException;
+import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
+
+import java.io.DataOutputStream;
+import java.io.IOException;
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Objects;
+
+public class AlterLogicalViewNode extends WritePlanNode {
+
+ /**
+ * A map from target path to source expression. Yht target path is the name
of this logical view,
+ * and the source expression is the data source of this view.
+ */
+ private Map<PartialPath, ViewExpression> viewPathToSourceMap;
+
+ /**
+ * This variable will be set in function splitByPartition() according to
analysis. And it will be
+ * set when creating new split nodes.
+ */
+ private TRegionReplicaSet regionReplicaSet = null;
+
+ public AlterLogicalViewNode(PlanNodeId id, Map<PartialPath, ViewExpression>
viewPathToSourceMap) {
+ super(id);
+ this.viewPathToSourceMap = viewPathToSourceMap;
+ }
+
+ public Map<PartialPath, ViewExpression> getViewPathToSourceMap() {
+ return viewPathToSourceMap;
+ }
+
+ // region Interfaces in WritePlanNode or PlanNode
+
+ @Override
+ public <R, C> R accept(PlanVisitor<R, C> visitor, C context) {
+ return visitor.visitAlterLogicalView(this, context);
+ }
+
+ @Override
+ public TRegionReplicaSet getRegionReplicaSet() {
+ return this.regionReplicaSet;
+ }
+
+ @Override
+ public List<PlanNode> getChildren() {
+ return new ArrayList<>();
+ }
+
+ @Override
+ public void addChild(PlanNode child) {
+ // do nothing. this node should never have any child
+ }
+
+ @Override
+ public PlanNode clone() {
+ // TODO: CRTODO, complete this method
+ throw new NotImplementedException("Clone of AlterLogicalNode is not
implemented");
+ }
+
+ @Override
+ public boolean equals(Object obj) {
+ if (this == obj) {
+ return true;
+ }
+ if (obj == null || getClass() != obj.getClass()) {
+ return false;
+ }
+ AlterLogicalViewNode that = (AlterLogicalViewNode) obj;
+ return (this.getPlanNodeId().equals(that.getPlanNodeId())
+ && Objects.equals(this.viewPathToSourceMap, that.viewPathToSourceMap));
+ }
+
+ @Override
+ public int allowedChildCount() {
+ // this node should never have any child
+ return NO_CHILD_ALLOWED;
+ }
+
+ @Override
+ public List<String> getOutputColumnNames() {
+ // TODO: CRTODO, complete this method
+ throw new NotImplementedException(
+ "getOutputColumnNames of AlterLogicalViewNode is not implemented");
+ }
+
+ @Override
+ protected void serializeAttributes(ByteBuffer byteBuffer) {
+ PlanNodeType.ALTER_LOGICAL_VIEW.serialize(byteBuffer);
+ // serialize other member variables for this node
+ ReadWriteIOUtils.write(this.viewPathToSourceMap.size(), byteBuffer);
+ for (Map.Entry<PartialPath, ViewExpression> entry :
viewPathToSourceMap.entrySet()) {
+ entry.getKey().serialize(byteBuffer);
+ ViewExpression.serialize(entry.getValue(), byteBuffer);
+ }
+ }
+
+ @Override
+ protected void serializeAttributes(DataOutputStream stream) throws
IOException {
+ PlanNodeType.ALTER_LOGICAL_VIEW.serialize(stream);
+ // serialize other member variables for this node
+ ReadWriteIOUtils.write(this.viewPathToSourceMap.size(), stream);
+ for (Map.Entry<PartialPath, ViewExpression> entry :
viewPathToSourceMap.entrySet()) {
+ entry.getKey().serialize(stream);
+ ViewExpression.serialize(entry.getValue(), stream);
+ }
+ }
+
+ public static AlterLogicalViewNode deserialize(ByteBuffer byteBuffer) {
+ // deserialize member variables
+ Map<PartialPath, ViewExpression> viewPathToSourceMap = new HashMap<>();
+ int size = byteBuffer.getInt();
+ PartialPath path;
+ ViewExpression viewExpression;
+ for (int i = 0; i < size; i++) {
+ path = (PartialPath) PathDeserializeUtil.deserialize(byteBuffer);
+ viewExpression = ViewExpression.deserialize(byteBuffer);
+ viewPathToSourceMap.put(path, viewExpression);
+ }
+ // deserialize PlanNodeId next
+ PlanNodeId planNodeId = PlanNodeId.deserialize(byteBuffer);
+ return new AlterLogicalViewNode(planNodeId, viewPathToSourceMap);
+ }
+
+ @Override
+ public List<WritePlanNode> splitByPartition(Analysis analysis) {
+ Map<TRegionReplicaSet, Map<PartialPath, ViewExpression>> splitMap = new
HashMap<>();
+ for (Map.Entry<PartialPath, ViewExpression> entry :
this.viewPathToSourceMap.entrySet()) {
+ // for each entry in the map for target path to source expression,
+ // build a map from TRegionReplicaSet to this entry.
+ // Please note that getSchemaRegionReplicaSet needs a device path as
parameter.
+ TRegionReplicaSet regionReplicaSet =
+
analysis.getSchemaPartitionInfo().getSchemaRegionReplicaSet(entry.getKey().getDevice());
+
+ // create a map if the key(regionReplicaSet) is not exists,
+ // then put this entry into this map(from regionReplicaSet to this entry)
+ splitMap
+ .computeIfAbsent(regionReplicaSet, k -> new HashMap<>())
+ .put(entry.getKey(), entry.getValue());
+ }
+
+ // split this node into several nodes according to their regionReplicaSet
+ List<WritePlanNode> result = new ArrayList<>();
+ for (Map.Entry<TRegionReplicaSet, Map<PartialPath, ViewExpression>> entry :
+ splitMap.entrySet()) {
+ // for each entry in splitMap, create a plan node.
+ result.add(new CreateLogicalViewNode(getPlanNodeId(), entry.getValue(),
entry.getKey()));
+ }
+ return result;
+ }
+ // endregion
+}
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/AlterLogicalViewStatement.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/AlterLogicalViewStatement.java
index fc7507e2b45..899d4f0b3be 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/AlterLogicalViewStatement.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/AlterLogicalViewStatement.java
@@ -27,7 +27,6 @@ import
org.apache.iotdb.db.mpp.plan.statement.IConfigStatement;
import org.apache.iotdb.db.mpp.plan.statement.Statement;
import org.apache.iotdb.db.mpp.plan.statement.StatementType;
import org.apache.iotdb.db.mpp.plan.statement.StatementVisitor;
-import org.apache.iotdb.db.mpp.plan.statement.component.IntoItem;
import org.apache.iotdb.db.mpp.plan.statement.crud.QueryStatement;
import java.util.List;
@@ -40,7 +39,6 @@ public class AlterLogicalViewStatement extends Statement
implements IConfigState
// the paths of sources
private ViewPaths sourcePaths;
private QueryStatement queryStatement;
- private IntoItem intoItem;
public AlterLogicalViewStatement() {
super();
@@ -103,15 +101,6 @@ public class AlterLogicalViewStatement extends Statement
implements IConfigState
this.targetPaths.setSuffixOfPathsGroup(suffixPaths);
this.targetPaths.generateFullPathsFromPathsGroup();
}
-
- public void setTargetIntoItem(IntoItem intoItem) {
- this.targetPaths.setViewPathType(ViewPathType.BATCH_GENERATION);
- this.intoItem = intoItem;
- }
-
- public IntoItem getIntoItem() {
- return this.intoItem;
- }
// endregion
@Override
diff --git
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/CreateLogicalViewStatement.java
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/CreateLogicalViewStatement.java
index 9ecd71fce8d..bd6a20d96ae 100644
---
a/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/CreateLogicalViewStatement.java
+++
b/server/src/main/java/org/apache/iotdb/db/mpp/plan/statement/metadata/view/CreateLogicalViewStatement.java
@@ -115,6 +115,10 @@ public class CreateLogicalViewStatement extends Statement {
this.queryStatement = queryStatement;
}
+ public void setQueryStatement(QueryStatement queryStatement) {
+ this.queryStatement = queryStatement;
+ }
+
/**
* This function must be called after analyzing query statement. Expressions
that analyzed should
* be set through here.
diff --git
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
index ba86131b29c..b31fd026a11 100644
---
a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++
b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@ -39,9 +39,11 @@ import org.apache.iotdb.commons.consensus.SchemaRegionId;
import org.apache.iotdb.commons.exception.IllegalPathException;
import org.apache.iotdb.commons.exception.MetadataException;
import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.path.PathDeserializeUtil;
import org.apache.iotdb.commons.path.PathPatternTree;
import org.apache.iotdb.commons.pipe.plugin.meta.PipePluginMeta;
import org.apache.iotdb.commons.pipe.task.meta.PipeMeta;
+import org.apache.iotdb.commons.schema.view.viewExpression.ViewExpression;
import org.apache.iotdb.commons.service.metric.MetricService;
import org.apache.iotdb.commons.service.metric.enums.Metric;
import org.apache.iotdb.commons.service.metric.enums.Tag;
@@ -103,6 +105,7 @@ import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.DeleteTimeS
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.PreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackPreDeactivateTemplateNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.RollbackSchemaBlackListNode;
+import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.AlterLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.ConstructLogicalViewBlackListNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.DeleteLogicalViewNode;
import
org.apache.iotdb.db.mpp.plan.planner.plan.node.metedata.write.view.RollbackLogicalViewBlackListNode;
@@ -127,6 +130,7 @@ import org.apache.iotdb.metrics.type.AutoGauge;
import org.apache.iotdb.metrics.utils.MetricLevel;
import org.apache.iotdb.mpp.rpc.thrift.IDataNodeRPCService;
import org.apache.iotdb.mpp.rpc.thrift.TActiveTriggerInstanceReq;
+import org.apache.iotdb.mpp.rpc.thrift.TAlterViewReq;
import org.apache.iotdb.mpp.rpc.thrift.TCancelFragmentInstanceReq;
import org.apache.iotdb.mpp.rpc.thrift.TCancelPlanFragmentReq;
import org.apache.iotdb.mpp.rpc.thrift.TCancelQueryReq;
@@ -194,6 +198,7 @@ import org.apache.iotdb.trigger.api.enums.TriggerEvent;
import org.apache.iotdb.tsfile.exception.NotImplementedException;
import org.apache.iotdb.tsfile.file.metadata.enums.TSDataType;
import org.apache.iotdb.tsfile.read.common.block.TsBlock;
+import org.apache.iotdb.tsfile.utils.ReadWriteIOUtils;
import org.apache.iotdb.tsfile.write.record.Tablet;
import com.google.common.collect.ImmutableList;
@@ -893,6 +898,36 @@ public class DataNodeInternalRPCServiceImpl implements
IDataNodeRPCService.Iface
});
}
+ @Override
+ public TSStatus alterView(TAlterViewReq req) throws TException {
+ List<TConsensusGroupId> consensusGroupIdList = req.getSchemaRegionIdList();
+ List<ByteBuffer> viewBinaryList = req.getViewBinaryList();
+ Map<TConsensusGroupId, Map<PartialPath, ViewExpression>>
schemaRegionRequestMap =
+ new HashMap<>();
+ for (int i = 0; i < consensusGroupIdList.size(); i++) {
+ ByteBuffer byteBuffer = viewBinaryList.get(i);
+ int size = ReadWriteIOUtils.readInt(byteBuffer);
+ Map<PartialPath, ViewExpression> viewMap = new HashMap<>();
+ for (int j = 0; j < size; j++) {
+ viewMap.put(
+ (PartialPath) PathDeserializeUtil.deserialize(byteBuffer),
+ ViewExpression.deserialize(byteBuffer));
+ }
+ schemaRegionRequestMap.put(consensusGroupIdList.get(i), viewMap);
+ }
+ return executeInternalSchemaTask(
+ consensusGroupIdList,
+ consensusGroupId -> {
+ RegionWriteExecutor executor = new RegionWriteExecutor();
+ return executor
+ .execute(
+ new SchemaRegionId(consensusGroupId.getId()),
+ new AlterLogicalViewNode(
+ new PlanNodeId(""),
schemaRegionRequestMap.get(consensusGroupId)))
+ .getStatus();
+ });
+ }
+
@Override
public TSStatus pushPipeMeta(TPushPipeMetaReq req) {
final List<PipeMeta> pipeMetas = new ArrayList<>();