This is an automated email from the ASF dual-hosted git repository.

healchow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 7fd2b327a [INLONG-4247][Manager] Add APIs for 
create/suspend/restart/delete stream (#4253)
7fd2b327a is described below

commit 7fd2b327ab4a35200f3298840e8676669e121c2e
Author: kipshi <[email protected]>
AuthorDate: Thu May 19 10:23:36 2022 +0800

    [INLONG-4247][Manager] Add APIs for create/suspend/restart/delete stream 
(#4253)
---
 .../inlong/manager/client/api/InlongGroupConf.java |   6 +-
 .../inlong/manager/common/enums/StreamStatus.java  |   9 +
 .../manager/service/core/ConsumptionService.java   |  10 -
 .../service/core/impl/ConsumptionServiceImpl.java  |  38 ---
 .../core/impl/InlongStreamProcessOperation.java    |  60 -----
 .../core/impl/WorkflowApproverServiceImpl.java     |   8 +-
 .../operation/ConsumptionProcessOperation.java     |  75 ++++++
 .../InlongGroupProcessOperation.java               |   4 +-
 .../operation/InlongStreamProcessOperation.java    | 275 +++++++++++++++++++++
 .../workflow/WorkflowDefinitionRegister.java       |  49 ++++
 .../service/workflow/WorkflowEngineConfig.java     |   4 +-
 .../service/workflow/WorkflowServiceImpl.java      |  25 +-
 .../core/impl/InlongGroupProcessOperationTest.java |   1 +
 .../src/main/resources/application-test.properties |   2 +-
 .../web/controller/ConsumptionController.java      |   7 +-
 .../web/controller/InlongGroupController.java      |   2 +-
 .../web/controller/InlongStreamController.java     |  52 ++++
 .../src/main/resources/application.properties      |   1 -
 .../manager/workflow/core/WorkflowEngine.java      |   7 -
 .../workflow/core/impl/WorkflowEngineImpl.java     |  12 +-
 20 files changed, 484 insertions(+), 163 deletions(-)

diff --git 
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/InlongGroupConf.java
 
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/InlongGroupConf.java
index 2664d1bbe..315d75cb2 100644
--- 
a/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/InlongGroupConf.java
+++ 
b/inlong-manager/manager-client/src/main/java/org/apache/inlong/manager/client/api/InlongGroupConf.java
@@ -57,13 +57,13 @@ public class InlongGroupConf {
     private String operator = "admin";
 
     @ApiModelProperty(value = "Whether to enable zookeeper? 0: disable, 1: 
enable")
-    private Integer enableZookeeper;
+    private Integer enableZookeeper = 1;
 
     @ApiModelProperty(value = "Whether to enable zookeeper? 0: disable, 1: 
enable")
-    private Integer enableCreateResource;
+    private Integer enableCreateResource = 1;
 
     @ApiModelProperty(value = "Whether to use lightweight mode, 0: false, 1: 
true")
-    private Integer lightweight;
+    private Integer lightweight = 0;
 
     @ApiModelProperty("Inlong cluster tag")
     private String inlongClusterTag;
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/StreamStatus.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/StreamStatus.java
index 8f1a51ed8..151922eee 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/StreamStatus.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/StreamStatus.java
@@ -55,4 +55,13 @@ public enum StreamStatus {
         return description;
     }
 
+    public static StreamStatus forCode(int code) {
+        for (StreamStatus status : values()) {
+            if (status.getCode() == code) {
+                return status;
+            }
+        }
+        throw new IllegalStateException(String.format("Illegal code=%s for 
StreamStatus", code));
+    }
+
 }
\ No newline at end of file
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
index bdad67949..322b1befa 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
@@ -23,7 +23,6 @@ import 
org.apache.inlong.manager.common.pojo.consumption.ConsumptionListVo;
 import org.apache.inlong.manager.common.pojo.consumption.ConsumptionQuery;
 import org.apache.inlong.manager.common.pojo.consumption.ConsumptionSummary;
 import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
-import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
 
 /**
  * Data consumption interface
@@ -88,15 +87,6 @@ public interface ConsumptionService {
      */
     Boolean delete(Integer id, String operator);
 
-    /**
-     * Start the application process
-     *
-     * @param id Data consumption id
-     * @param operator Operator
-     * @return WorkflowProcess information
-     */
-    WorkflowResult startProcess(Integer id, String operator);
-
     /**
      * Save the consumer group info for Sort to the database
      */
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
index 6141be364..4b632d696 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
@@ -39,8 +39,6 @@ import 
org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
 import org.apache.inlong.manager.common.pojo.group.InlongGroupTopicResponse;
 import org.apache.inlong.manager.common.pojo.stream.InlongStreamTopicResponse;
 import org.apache.inlong.manager.common.pojo.user.UserRoleCode;
-import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
-import 
org.apache.inlong.manager.common.pojo.workflow.form.NewConsumptionProcessForm;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.LoginUserUtils;
 import org.apache.inlong.manager.common.util.Preconditions;
@@ -53,8 +51,6 @@ import 
org.apache.inlong.manager.dao.mapper.InlongGroupEntityMapper;
 import org.apache.inlong.manager.service.core.ConsumptionService;
 import org.apache.inlong.manager.service.core.InlongGroupService;
 import org.apache.inlong.manager.service.core.InlongStreamService;
-import org.apache.inlong.manager.service.workflow.ProcessName;
-import org.apache.inlong.manager.service.workflow.WorkflowService;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.stereotype.Service;
 import org.springframework.transaction.annotation.Transactional;
@@ -86,8 +82,6 @@ public class ConsumptionServiceImpl implements 
ConsumptionService {
     @Autowired
     private ConsumptionPulsarEntityMapper consumptionPulsarMapper;
     @Autowired
-    private WorkflowService workflowService;
-    @Autowired
     private InlongGroupService groupService;
     @Autowired
     private InlongStreamService streamService;
@@ -323,24 +317,6 @@ public class ConsumptionServiceImpl implements 
ConsumptionService {
         return true;
     }
 
-    @Override
-    public WorkflowResult startProcess(Integer id, String operation) {
-        ConsumptionInfo consumptionInfo = this.get(id);
-        
Preconditions.checkTrue(ConsumptionStatus.ALLOW_START_WORKFLOW_STATUS.contains(
-                        
ConsumptionStatus.fromStatus(consumptionInfo.getStatus())),
-                "current status not allow start workflow");
-
-        ConsumptionEntity updateConsumptionEntity = new ConsumptionEntity();
-        updateConsumptionEntity.setId(consumptionInfo.getId());
-        updateConsumptionEntity.setModifyTime(new Date());
-        
updateConsumptionEntity.setStatus(ConsumptionStatus.WAIT_APPROVE.getStatus());
-        int success = 
this.consumptionMapper.updateByPrimaryKeySelective(updateConsumptionEntity);
-        Preconditions.checkTrue(success == 1, "update consumption failed");
-
-        return workflowService.start(ProcessName.NEW_CONSUMPTION_PROCESS, 
operation,
-                genNewConsumptionProcessForm(consumptionInfo));
-    }
-
     @Override
     public void saveSortConsumption(InlongGroupInfo groupInfo, String topic, 
String consumerGroup) {
         String groupId = groupInfo.getInlongGroupId();
@@ -380,20 +356,6 @@ public class ConsumptionServiceImpl implements 
ConsumptionService {
         log.debug("success save consumption, groupId={}, topic={}, consumer 
group={}", groupId, topic, consumerGroup);
     }
 
-    private NewConsumptionProcessForm 
genNewConsumptionProcessForm(ConsumptionInfo consumptionInfo) {
-        NewConsumptionProcessForm form = new NewConsumptionProcessForm();
-        Integer id = consumptionInfo.getId();
-        MQType mqType = MQType.forType(consumptionInfo.getMqType());
-        if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
-            ConsumptionPulsarEntity consumptionPulsarEntity = 
consumptionPulsarMapper.selectByConsumptionId(id);
-            ConsumptionPulsarInfo pulsarInfo = 
CommonBeanUtils.copyProperties(consumptionPulsarEntity,
-                    ConsumptionPulsarInfo::new);
-            consumptionInfo.setMqExtInfo(pulsarInfo);
-        }
-        form.setConsumptionInfo(consumptionInfo);
-        return form;
-    }
-
     private ConsumptionEntity saveConsumption(ConsumptionInfo info, String 
operator, Date now) {
         ConsumptionEntity entity = CommonBeanUtils.copyProperties(info, 
ConsumptionEntity::new);
         entity.setStatus(ConsumptionStatus.WAIT_ASSIGN.getStatus());
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/InlongStreamProcessOperation.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/InlongStreamProcessOperation.java
deleted file mode 100644
index 9cfda591a..000000000
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/InlongStreamProcessOperation.java
+++ /dev/null
@@ -1,60 +0,0 @@
-/*
- * 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.inlong.manager.service.core.impl;
-
-import com.google.common.util.concurrent.ThreadFactoryBuilder;
-import lombok.extern.slf4j.Slf4j;
-import org.apache.inlong.manager.service.core.InlongGroupService;
-import org.apache.inlong.manager.service.core.InlongStreamService;
-import org.apache.inlong.manager.service.workflow.WorkflowService;
-import org.springframework.beans.factory.annotation.Autowired;
-import org.springframework.stereotype.Service;
-
-import java.util.concurrent.ExecutorService;
-import java.util.concurrent.LinkedBlockingQueue;
-import java.util.concurrent.ThreadPoolExecutor;
-import java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy;
-import java.util.concurrent.TimeUnit;
-
-/**
- * Operation related to inlong stream process
- */
-@Service
-@Slf4j
-public class InlongStreamProcessOperation {
-
-    private final ExecutorService executorService = new ThreadPoolExecutor(
-            20,
-            40,
-            0L,
-            TimeUnit.MILLISECONDS,
-            new LinkedBlockingQueue<>(),
-            new 
ThreadFactoryBuilder().setNameFormat("inlong-stream-process-%s").build(),
-            new CallerRunsPolicy());
-
-    @Autowired
-    private InlongGroupService groupService;
-
-    @Autowired
-    private InlongStreamService streamService;
-
-    @Autowired
-    private WorkflowService workflowService;
-
-
-}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/WorkflowApproverServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/WorkflowApproverServiceImpl.java
index 60f78a127..20cb9ebb3 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/WorkflowApproverServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/WorkflowApproverServiceImpl.java
@@ -28,7 +28,7 @@ import org.apache.inlong.manager.common.util.Preconditions;
 import org.apache.inlong.manager.dao.entity.WorkflowApproverEntity;
 import org.apache.inlong.manager.dao.mapper.WorkflowApproverEntityMapper;
 import org.apache.inlong.manager.service.core.WorkflowApproverService;
-import org.apache.inlong.manager.workflow.core.WorkflowEngine;
+import org.apache.inlong.manager.workflow.core.ProcessDefinitionService;
 import org.apache.inlong.manager.workflow.definition.UserTask;
 import org.apache.inlong.manager.workflow.definition.WorkflowProcess;
 import org.apache.inlong.manager.workflow.definition.WorkflowTask;
@@ -52,7 +52,7 @@ public class WorkflowApproverServiceImpl implements 
WorkflowApproverService {
     @Autowired
     private WorkflowApproverEntityMapper workflowApproverMapper;
     @Autowired
-    private WorkflowEngine workflowEngine;
+    private ProcessDefinitionService processDefinitionService;
 
     @Override
     public List<String> getApprovers(String processName, String taskName, 
WorkflowApproverFilterContext context) {
@@ -83,7 +83,7 @@ public class WorkflowApproverServiceImpl implements 
WorkflowApproverService {
         List<WorkflowApproverEntity> entityList = 
workflowApproverMapper.selectByQuery(query);
         List<WorkflowApprover> approverList = 
CommonBeanUtils.copyListProperties(entityList, WorkflowApprover::new);
         approverList.forEach(config -> {
-            WorkflowProcess process = 
workflowEngine.processDefinitionService().getByName(config.getProcessName());
+            WorkflowProcess process = 
processDefinitionService.getByName(config.getProcessName());
             if (process != null) {
                 config.setProcessDisplayName(process.getDisplayName());
                 
config.setTaskDisplayName(Optional.ofNullable(process.getTaskByName(config.getTaskName())).map(
@@ -102,7 +102,7 @@ public class WorkflowApproverServiceImpl implements 
WorkflowApproverService {
         approver.setModifier(operator);
         approver.setCreator(operator);
 
-        WorkflowProcess process = 
workflowEngine.processDefinitionService().getByName(approver.getProcessName());
+        WorkflowProcess process = 
processDefinitionService.getByName(approver.getProcessName());
         Preconditions.checkNotNull(process, "process not exit with name: " + 
approver.getProcessName());
         WorkflowTask task = process.getTaskByName(approver.getTaskName());
         Preconditions.checkNotNull(task, "task not exit with name: " + 
approver.getTaskName());
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/ConsumptionProcessOperation.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/ConsumptionProcessOperation.java
new file mode 100644
index 000000000..5a18360f5
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/ConsumptionProcessOperation.java
@@ -0,0 +1,75 @@
+/*
+ * 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.inlong.manager.service.core.operation;
+
+import lombok.extern.slf4j.Slf4j;
+import org.apache.inlong.manager.common.enums.ConsumptionStatus;
+import org.apache.inlong.manager.common.enums.MQType;
+import org.apache.inlong.manager.common.pojo.consumption.ConsumptionInfo;
+import org.apache.inlong.manager.common.pojo.consumption.ConsumptionPulsarInfo;
+import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
+import 
org.apache.inlong.manager.common.pojo.workflow.form.NewConsumptionProcessForm;
+import org.apache.inlong.manager.common.util.CommonBeanUtils;
+import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.dao.entity.ConsumptionPulsarEntity;
+import org.apache.inlong.manager.dao.mapper.ConsumptionPulsarEntityMapper;
+import org.apache.inlong.manager.service.core.ConsumptionService;
+import org.apache.inlong.manager.service.workflow.ProcessName;
+import org.apache.inlong.manager.service.workflow.WorkflowService;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+@Service
+@Slf4j
+public class ConsumptionProcessOperation {
+
+    @Autowired
+    private ConsumptionService consumptionService;
+    @Autowired
+    private WorkflowService workflowService;
+    @Autowired
+    private ConsumptionPulsarEntityMapper consumptionPulsarMapper;
+
+    public WorkflowResult startProcess(Integer id, String operator) {
+        ConsumptionInfo consumptionInfo = consumptionService.get(id);
+        
Preconditions.checkTrue(ConsumptionStatus.ALLOW_START_WORKFLOW_STATUS.contains(
+                        
ConsumptionStatus.fromStatus(consumptionInfo.getStatus())),
+                "current status not allow start workflow");
+
+        consumptionInfo.setStatus(ConsumptionStatus.WAIT_APPROVE.getStatus());
+        boolean isSuccess = 
consumptionService.update(consumptionInfo,operator);
+        Preconditions.checkTrue(isSuccess, "update consumption failed");
+
+        return workflowService.start(ProcessName.NEW_CONSUMPTION_PROCESS, 
operator,
+                genNewConsumptionProcessForm(consumptionInfo));
+    }
+
+    private NewConsumptionProcessForm 
genNewConsumptionProcessForm(ConsumptionInfo consumptionInfo) {
+        NewConsumptionProcessForm form = new NewConsumptionProcessForm();
+        Integer id = consumptionInfo.getId();
+        MQType mqType = MQType.forType(consumptionInfo.getMqType());
+        if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
+            ConsumptionPulsarEntity consumptionPulsarEntity = 
consumptionPulsarMapper.selectByConsumptionId(id);
+            ConsumptionPulsarInfo pulsarInfo = 
CommonBeanUtils.copyProperties(consumptionPulsarEntity,
+                    ConsumptionPulsarInfo::new);
+            consumptionInfo.setMqExtInfo(pulsarInfo);
+        }
+        form.setConsumptionInfo(consumptionInfo);
+        return form;
+    }
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperation.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongGroupProcessOperation.java
similarity index 98%
rename from 
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperation.java
rename to 
inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongGroupProcessOperation.java
index a7431369b..a429dafc4 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperation.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongGroupProcessOperation.java
@@ -15,7 +15,7 @@
  * limitations under the License.
  */
 
-package org.apache.inlong.manager.service.core.impl;
+package org.apache.inlong.manager.service.core.operation;
 
 import com.google.common.util.concurrent.ThreadFactoryBuilder;
 import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
@@ -247,7 +247,7 @@ public class InlongGroupProcessOperation {
     /**
      * Generate the form of [New Group Workflow]
      */
-    public NewGroupProcessForm genNewGroupProcessForm(String groupId) {
+    private NewGroupProcessForm genNewGroupProcessForm(String groupId) {
         NewGroupProcessForm form = new NewGroupProcessForm();
         InlongGroupInfo groupInfo = groupService.get(groupId);
         form.setGroupInfo(groupInfo);
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongStreamProcessOperation.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongStreamProcessOperation.java
new file mode 100644
index 000000000..eb6494161
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/operation/InlongStreamProcessOperation.java
@@ -0,0 +1,275 @@
+/*
+ * 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.inlong.manager.service.core.operation;
+
+import com.google.common.util.concurrent.ThreadFactoryBuilder;
+import lombok.extern.slf4j.Slf4j;
+import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
+import org.apache.inlong.manager.common.enums.GroupOperateType;
+import org.apache.inlong.manager.common.enums.GroupStatus;
+import org.apache.inlong.manager.common.enums.ProcessStatus;
+import org.apache.inlong.manager.common.enums.StreamStatus;
+import org.apache.inlong.manager.common.exceptions.BusinessException;
+import org.apache.inlong.manager.common.pojo.group.InlongGroupInfo;
+import org.apache.inlong.manager.common.pojo.stream.InlongStreamInfo;
+import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
+import 
org.apache.inlong.manager.common.pojo.workflow.form.StreamResourceProcessForm;
+import org.apache.inlong.manager.service.core.InlongGroupService;
+import org.apache.inlong.manager.service.core.InlongStreamService;
+import org.apache.inlong.manager.service.workflow.ProcessName;
+import org.apache.inlong.manager.service.workflow.WorkflowService;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.ThreadPoolExecutor.CallerRunsPolicy;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Operation related to inlong stream process
+ */
+@Service
+@Slf4j
+public class InlongStreamProcessOperation {
+
+    private final ExecutorService executorService = new ThreadPoolExecutor(
+            20,
+            40,
+            0L,
+            TimeUnit.MILLISECONDS,
+            new LinkedBlockingQueue<>(),
+            new 
ThreadFactoryBuilder().setNameFormat("inlong-stream-process-%s").build(),
+            new CallerRunsPolicy());
+
+    @Autowired
+    private InlongGroupService groupService;
+
+    @Autowired
+    private InlongStreamService streamService;
+
+    @Autowired
+    private WorkflowService workflowService;
+
+    /**
+     * Create stream in synchronous/asynchronous way.
+     *
+     * @return
+     */
+    public boolean startProcess(String groupId, String streamId, String 
operator, boolean sync) {
+        log.info("StartProcess for groupId={}, streamId={}", groupId, 
streamId);
+        InlongGroupInfo groupInfo = groupService.get(groupId);
+        if (groupInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.GROUP_NOT_FOUND);
+        }
+        GroupStatus groupStatus = GroupStatus.forCode(groupInfo.getStatus());
+        if (groupStatus != GroupStatus.CONFIG_SUCCESSFUL && groupStatus != 
GroupStatus.RESTARTED) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, status=%s not correct for 
stream start", groupId, groupStatus));
+        }
+        InlongStreamInfo streamInfo = streamService.get(groupId, streamId);
+        if (streamInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.STREAM_NOT_FOUND);
+        }
+        StreamStatus status = StreamStatus.forCode(streamInfo.getStatus());
+        if (status == StreamStatus.CONFIG_ING || status == 
StreamStatus.CONFIG_SUCCESSFUL) {
+            log.warn("GroupId={}, StreamId={} is already in {}", groupId, 
streamId, status);
+            return true;
+        }
+        if (status != StreamStatus.NEW || status != 
StreamStatus.CONFIG_FAILED) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, StreamId=%s, status=%s not 
correct for stream start", groupId, streamId,
+                            status));
+        }
+        StreamResourceProcessForm processForm = 
genStreamProcessForm(groupInfo, streamInfo, GroupOperateType.INIT);
+        ProcessName processName = ProcessName.CREATE_STREAM_RESOURCE;
+        if (sync) {
+            WorkflowResult workflowResult = workflowService.start(processName, 
operator,
+                    processForm);
+            ProcessStatus processStatus = 
workflowResult.getProcessInfo().getStatus();
+            return processStatus == ProcessStatus.COMPLETED;
+        } else {
+            executorService.execute(
+                    () -> workflowService.start(processName, operator, 
processForm));
+            return true;
+        }
+    }
+
+    /**
+     * Suspend stream in synchronous/asynchronous way.
+     *
+     * @return
+     */
+    public boolean suspendProcess(String groupId, String streamId, String 
operator, boolean sync) {
+        log.info("SuspendProcess for groupId={}, streamId={}", groupId, 
streamId);
+        InlongGroupInfo groupInfo = groupService.get(groupId);
+        if (groupInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.GROUP_NOT_FOUND);
+        }
+        GroupStatus groupStatus = GroupStatus.forCode(groupInfo.getStatus());
+        if (groupStatus != GroupStatus.CONFIG_SUCCESSFUL
+                && groupStatus != GroupStatus.RESTARTED
+                && groupStatus != GroupStatus.SUSPENDED) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, status=%s not correct for 
stream suspend", groupId, groupStatus));
+        }
+        InlongStreamInfo streamInfo = streamService.get(groupId, streamId);
+        if (streamInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.STREAM_NOT_FOUND);
+        }
+        StreamStatus status = StreamStatus.forCode(streamInfo.getStatus());
+        if (status == StreamStatus.SUSPENDED || status == 
StreamStatus.SUSPENDING) {
+            log.warn("GroupId={}, StreamId={} is already in {}", groupId, 
streamId, status);
+            return true;
+        }
+        if (status != StreamStatus.CONFIG_SUCCESSFUL && status != 
StreamStatus.RESTARTED) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, StreamId=%s, status=%s not 
correct for stream suspend", groupId,
+                            streamId,
+                            status));
+        }
+        StreamResourceProcessForm processForm = 
genStreamProcessForm(groupInfo, streamInfo, GroupOperateType.SUSPEND);
+        ProcessName processName = ProcessName.SUSPEND_STREAM_RESOURCE;
+        if (sync) {
+            WorkflowResult workflowResult = workflowService.start(processName, 
operator,
+                    processForm);
+            ProcessStatus processStatus = 
workflowResult.getProcessInfo().getStatus();
+            return processStatus == ProcessStatus.COMPLETED;
+        } else {
+            executorService.execute(
+                    () -> workflowService.start(processName, operator, 
processForm));
+            return true;
+        }
+    }
+
+    /**
+     * Restart stream in synchronous/asynchronous way.
+     *
+     * @return
+     */
+    public boolean restartProcess(String groupId, String streamId, String 
operator, boolean sync) {
+        log.info("RestartProcess for groupId={}, streamId={}", groupId, 
streamId);
+        InlongGroupInfo groupInfo = groupService.get(groupId);
+        if (groupInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.GROUP_NOT_FOUND);
+        }
+        GroupStatus groupStatus = GroupStatus.forCode(groupInfo.getStatus());
+        if (groupStatus != GroupStatus.CONFIG_SUCCESSFUL
+                && groupStatus != GroupStatus.RESTARTED) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, status=%s not correct for 
stream restart", groupId, groupStatus));
+        }
+        InlongStreamInfo streamInfo = streamService.get(groupId, streamId);
+        if (streamInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.STREAM_NOT_FOUND);
+        }
+        StreamStatus status = StreamStatus.forCode(streamInfo.getStatus());
+        if (status == StreamStatus.RESTARTED || status == 
StreamStatus.RESTARTING) {
+            log.warn("GroupId={}, StreamId={} is already in {}", groupId, 
streamId, status);
+            return true;
+        }
+        if (status != StreamStatus.SUSPENDED) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, StreamId=%s, status=%s not 
correct for stream restart", groupId,
+                            streamId,
+                            status));
+        }
+        StreamResourceProcessForm processForm = 
genStreamProcessForm(groupInfo, streamInfo, GroupOperateType.RESTART);
+        ProcessName processName = ProcessName.RESTART_STREAM_RESOURCE;
+        if (sync) {
+            WorkflowResult workflowResult = workflowService.start(processName, 
operator,
+                    processForm);
+            ProcessStatus processStatus = 
workflowResult.getProcessInfo().getStatus();
+            return processStatus == ProcessStatus.COMPLETED;
+        } else {
+            executorService.execute(
+                    () -> workflowService.start(processName, operator, 
processForm));
+            return true;
+        }
+    }
+
+    /**
+     * Restart stream in synchronous/asynchronous way.
+     *
+     * @return
+     */
+    public boolean deleteProcess(String groupId, String streamId, String 
operator, boolean sync) {
+        log.info("DeleteProcess for groupId={}, streamId={}", groupId, 
streamId);
+        InlongGroupInfo groupInfo = groupService.get(groupId);
+        if (groupInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.GROUP_NOT_FOUND);
+        }
+        GroupStatus groupStatus = GroupStatus.forCode(groupInfo.getStatus());
+        if (groupStatus != GroupStatus.CONFIG_SUCCESSFUL
+                && groupStatus != GroupStatus.RESTARTED
+                && groupStatus != GroupStatus.SUSPENDED
+                && groupStatus != GroupStatus.DELETING) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, status=%s not correct for 
stream delete", groupId, groupStatus));
+        }
+        InlongStreamInfo streamInfo = streamService.get(groupId, streamId);
+        if (streamInfo == null) {
+            throw new BusinessException(ErrorCodeEnum.STREAM_NOT_FOUND);
+        }
+        StreamStatus status = StreamStatus.forCode(streamInfo.getStatus());
+        if (status == StreamStatus.DELETED || status == StreamStatus.DELETING) 
{
+            log.warn("GroupId={}, StreamId={} is already in {}", groupId, 
streamId, status);
+            return true;
+        }
+        if (status == StreamStatus.CONFIG_ING
+                || status == StreamStatus.RESTARTING
+                || status == StreamStatus.SUSPENDING) {
+            throw new BusinessException(
+                    String.format("GroupId=%s, StreamId=%s, status=%s not 
correct for stream delete", groupId,
+                            streamId,
+                            status));
+        }
+        StreamResourceProcessForm processForm = 
genStreamProcessForm(groupInfo, streamInfo, GroupOperateType.DELETE);
+        ProcessName processName = ProcessName.DELETE_STREAM_RESOURCE;
+        if (sync) {
+            WorkflowResult workflowResult = workflowService.start(processName, 
operator,
+                    processForm);
+            ProcessStatus processStatus = 
workflowResult.getProcessInfo().getStatus();
+            if (processStatus == ProcessStatus.COMPLETED) {
+                return streamService.delete(groupId, streamId, operator);
+            } else {
+                return false;
+            }
+        } else {
+            executorService.execute(
+                    () -> {
+                        WorkflowResult workflowResult = 
workflowService.start(processName, operator, processForm);
+                        ProcessStatus processStatus = 
workflowResult.getProcessInfo().getStatus();
+                        if (processStatus == ProcessStatus.COMPLETED) {
+                            streamService.delete(groupId, streamId, operator);
+                        }
+                    });
+            return true;
+        }
+    }
+
+    private StreamResourceProcessForm genStreamProcessForm(InlongGroupInfo 
groupInfo, InlongStreamInfo streamInfo,
+            GroupOperateType operateType) {
+        StreamResourceProcessForm processForm = new 
StreamResourceProcessForm();
+        processForm.setGroupInfo(groupInfo);
+        processForm.setStreamInfo(streamInfo);
+        processForm.setGroupOperateType(operateType);
+        return processForm;
+    }
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowDefinitionRegister.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowDefinitionRegister.java
new file mode 100644
index 000000000..7f963f1aa
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowDefinitionRegister.java
@@ -0,0 +1,49 @@
+/*
+ * 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.inlong.manager.service.workflow;
+
+import lombok.extern.slf4j.Slf4j;
+import org.apache.inlong.manager.workflow.core.WorkflowEngine;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import javax.annotation.PostConstruct;
+import java.util.List;
+
+@Service
+@Slf4j
+public class WorkflowDefinitionRegister {
+
+    @Autowired
+    private WorkflowEngine workflowEngine;
+    @Autowired
+    private List<WorkflowDefinition> workflowDefinitions;
+
+    @PostConstruct
+    public void registerDefinition() {
+        workflowDefinitions.forEach(definition -> {
+            try {
+                
workflowEngine.processDefinitionService().register(definition.defineProcess());
+                log.info("success register workflow definition: {}", 
definition.getProcessName());
+            } catch (Exception e) {
+                log.error("failed to register workflow definition {}", 
definition.getProcessName(), e);
+            }
+        });
+    }
+
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowEngineConfig.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowEngineConfig.java
index d74a00be6..bc122e500 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowEngineConfig.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowEngineConfig.java
@@ -17,6 +17,7 @@
 
 package org.apache.inlong.manager.service.workflow;
 
+import lombok.extern.slf4j.Slf4j;
 import org.apache.inlong.manager.dao.mapper.WorkflowEventLogEntityMapper;
 import org.apache.inlong.manager.dao.mapper.WorkflowProcessEntityMapper;
 import org.apache.inlong.manager.dao.mapper.WorkflowTaskEntityMapper;
@@ -38,6 +39,7 @@ import 
org.springframework.transaction.PlatformTransactionManager;
  */
 @Component
 @Configuration
+@Slf4j
 public class WorkflowEngineConfig {
 
     @Autowired
@@ -54,7 +56,7 @@ public class WorkflowEngineConfig {
     private PlatformTransactionManager platformTransactionManager;
 
     @Bean
-    public WorkflowEngine workflowEngineer() {
+    public WorkflowEngine workflowEngine() {
         WorkflowConfig workflowConfig = new WorkflowConfig()
                 .setQueryService(queryService)
                 .setProcessEntityMapper(processEntityMapper)
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowServiceImpl.java
index 9591c135d..b10cb13c2 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/WorkflowServiceImpl.java
@@ -56,7 +56,6 @@ import org.springframework.stereotype.Service;
 import org.springframework.transaction.annotation.Transactional;
 import org.springframework.util.CollectionUtils;
 
-import javax.annotation.PostConstruct;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
@@ -71,31 +70,11 @@ public class WorkflowServiceImpl implements WorkflowService 
{
 
     private static final Logger LOGGER = 
LoggerFactory.getLogger(WorkflowServiceImpl.class);
 
-    private final WorkflowEngine workflowEngine;
-
-    @Autowired
-    private WorkflowQueryService queryService;
     @Autowired
-    private List<WorkflowDefinition> workflowDefinitions;
+    private WorkflowEngine workflowEngine;
 
     @Autowired
-    public WorkflowServiceImpl(WorkflowEngine workflowEngine) {
-        this.workflowEngine = workflowEngine;
-    }
-
-    @PostConstruct
-    private void init() {
-        LOGGER.info("start init workflow service");
-        workflowDefinitions.forEach(definition -> {
-            try {
-                
workflowEngine.processDefinitionService().register(definition.defineProcess());
-                LOGGER.info("success register workflow definition: {}", 
definition.getProcessName());
-            } catch (Exception e) {
-                LOGGER.error("failed to register workflow definition {}", 
definition.getProcessName(), e);
-            }
-        });
-        LOGGER.info("success init workflow service");
-    }
+    private WorkflowQueryService queryService;
 
     @Override
     @Transactional(noRollbackFor = WorkflowNoRollbackException.class, 
rollbackFor = Exception.class)
diff --git 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
index 5287cb074..0401de732 100644
--- 
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
+++ 
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
@@ -27,6 +27,7 @@ import 
org.apache.inlong.manager.common.pojo.workflow.ProcessResponse;
 import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
 import org.apache.inlong.manager.service.ServiceBaseTest;
 import org.apache.inlong.manager.service.core.InlongGroupService;
+import 
org.apache.inlong.manager.service.core.operation.InlongGroupProcessOperation;
 import org.apache.inlong.manager.service.mocks.MockPlugin;
 import 
org.apache.inlong.manager.service.workflow.listener.GroupTaskListenerFactory;
 import org.junit.Assert;
diff --git 
a/inlong-manager/manager-test/src/main/resources/application-test.properties 
b/inlong-manager/manager-test/src/main/resources/application-test.properties
index 471175492..371c65340 100644
--- a/inlong-manager/manager-test/src/main/resources/application-test.properties
+++ b/inlong-manager/manager-test/src/main/resources/application-test.properties
@@ -100,4 +100,4 @@ common.http-client.validateAfterInactivity=5000
 common.http-client.connectionTimeout=3000
 common.http-client.readTimeout=10000
 common.http-client.connectionRequestTimeout=3000
-spring.main.allow-circular-references=true
\ No newline at end of file
+spring.main.allow-circular-references=false
\ No newline at end of file
diff --git 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/ConsumptionController.java
 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/ConsumptionController.java
index c26b9b9c6..b32965173 100644
--- 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/ConsumptionController.java
+++ 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/ConsumptionController.java
@@ -27,10 +27,11 @@ import 
org.apache.inlong.manager.common.pojo.consumption.ConsumptionInfo;
 import org.apache.inlong.manager.common.pojo.consumption.ConsumptionListVo;
 import org.apache.inlong.manager.common.pojo.consumption.ConsumptionQuery;
 import org.apache.inlong.manager.common.pojo.consumption.ConsumptionSummary;
+import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
 import org.apache.inlong.manager.common.util.LoginUserUtils;
 import org.apache.inlong.manager.service.core.ConsumptionService;
+import 
org.apache.inlong.manager.service.core.operation.ConsumptionProcessOperation;
 import org.apache.inlong.manager.service.core.operationlog.OperationLog;
-import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.validation.annotation.Validated;
 import org.springframework.web.bind.annotation.DeleteMapping;
@@ -51,6 +52,8 @@ public class ConsumptionController {
 
     @Autowired
     private ConsumptionService consumptionService;
+    @Autowired
+    private ConsumptionProcessOperation processOperation;
 
     @GetMapping("/summary")
     @ApiOperation(value = "Get data consumption summary")
@@ -106,7 +109,7 @@ public class ConsumptionController {
     @ApiImplicitParam(name = "id", value = "Consumption ID", dataTypeClass = 
Integer.class, required = true)
     public Response<WorkflowResult> startProcess(@PathVariable(name = "id") 
Integer id) {
         String username = LoginUserUtils.getLoginUserDetail().getUserName();
-        return Response.success(this.consumptionService.startProcess(id, 
username));
+        return Response.success(processOperation.startProcess(id, username));
     }
 
 }
diff --git 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongGroupController.java
 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongGroupController.java
index 6205063ec..a551da3cd 100644
--- 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongGroupController.java
+++ 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongGroupController.java
@@ -33,7 +33,7 @@ import 
org.apache.inlong.manager.common.pojo.group.InlongGroupTopicResponse;
 import org.apache.inlong.manager.common.pojo.workflow.WorkflowResult;
 import org.apache.inlong.manager.common.util.LoginUserUtils;
 import org.apache.inlong.manager.service.core.InlongGroupService;
-import org.apache.inlong.manager.service.core.impl.InlongGroupProcessOperation;
+import 
org.apache.inlong.manager.service.core.operation.InlongGroupProcessOperation;
 import org.apache.inlong.manager.service.core.operationlog.OperationLog;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.web.bind.annotation.PathVariable;
diff --git 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
index fd36fe350..be69bfa05 100644
--- 
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
+++ 
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongStreamController.java
@@ -34,6 +34,7 @@ import 
org.apache.inlong.manager.common.pojo.stream.StreamBriefResponse;
 import org.apache.inlong.manager.common.pojo.user.UserRoleCode;
 import org.apache.inlong.manager.common.util.LoginUserUtils;
 import org.apache.inlong.manager.service.core.InlongStreamService;
+import 
org.apache.inlong.manager.service.core.operation.InlongStreamProcessOperation;
 import org.apache.inlong.manager.service.core.operationlog.OperationLog;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.web.bind.annotation.PathVariable;
@@ -55,6 +56,8 @@ public class InlongStreamController {
 
     @Autowired
     private InlongStreamService streamService;
+    @Autowired
+    private InlongStreamProcessOperation streamProcessOperation;
 
     @RequestMapping(value = "/save", method = RequestMethod.POST)
     @OperationLog(operation = OperationType.CREATE)
@@ -123,6 +126,55 @@ public class InlongStreamController {
         return Response.success(streamService.update(request, username));
     }
 
+    @RequestMapping(value = "/startProcess/{groupId}/{streamId}", method = 
RequestMethod.POST)
+    @ApiOperation(value = "Start inlong stream")
+    @ApiImplicitParams({
+            @ApiImplicitParam(name = "groupId", dataTypeClass = String.class, 
required = true),
+            @ApiImplicitParam(name = "streamId", dataTypeClass = String.class, 
required = true)
+    })
+    public Response<Boolean> startProcess(@PathVariable String groupId, 
@PathVariable String streamId,
+            @RequestParam boolean sync) {
+        String operator = LoginUserUtils.getLoginUserDetail().getUserName();
+        return Response.success(streamProcessOperation.startProcess(groupId, 
streamId, operator, sync));
+    }
+
+    @RequestMapping(value = "/suspendProcess/{groupId}/{streamId}", method = 
RequestMethod.POST)
+    @ApiOperation(value = "Suspend inlong stream")
+    @ApiImplicitParams({
+            @ApiImplicitParam(name = "groupId", dataTypeClass = String.class, 
required = true),
+            @ApiImplicitParam(name = "streamId", dataTypeClass = String.class, 
required = true)
+    })
+    public Response<Boolean> suspendProcess(@PathVariable String groupId, 
@PathVariable String streamId,
+            @RequestParam boolean sync) {
+        String operator = LoginUserUtils.getLoginUserDetail().getUserName();
+        return Response.success(streamProcessOperation.suspendProcess(groupId, 
streamId, operator, sync));
+    }
+
+    @RequestMapping(value = "/restartProcess/{groupId}/{streamId}", method = 
RequestMethod.POST)
+    @ApiOperation(value = "Restart inlong stream")
+    @ApiImplicitParams({
+            @ApiImplicitParam(name = "groupId", dataTypeClass = String.class, 
required = true),
+            @ApiImplicitParam(name = "streamId", dataTypeClass = String.class, 
required = true)
+    })
+    public Response<Boolean> restartProcess(@PathVariable String groupId, 
@PathVariable String streamId,
+            @RequestParam boolean sync) {
+        String operator = LoginUserUtils.getLoginUserDetail().getUserName();
+        return Response.success(streamProcessOperation.restartProcess(groupId, 
streamId, operator, sync));
+    }
+
+    @RequestMapping(value = "/deleteProcess/{groupId}/{streamId}", method = 
RequestMethod.POST)
+    @ApiOperation(value = "Delete inlong stream")
+    @ApiImplicitParams({
+            @ApiImplicitParam(name = "groupId", dataTypeClass = String.class, 
required = true),
+            @ApiImplicitParam(name = "streamId", dataTypeClass = String.class, 
required = true)
+    })
+    public Response<Boolean> deleteProcess(@PathVariable String groupId, 
@PathVariable String streamId,
+            @RequestParam boolean sync) {
+        String operator = LoginUserUtils.getLoginUserDetail().getUserName();
+        return Response.success(streamProcessOperation.deleteProcess(groupId, 
streamId, operator, sync));
+    }
+
+    @Deprecated
     @RequestMapping(value = "/delete", method = RequestMethod.DELETE)
     @OperationLog(operation = OperationType.DELETE)
     @ApiOperation(value = "Delete inlong stream info")
diff --git 
a/inlong-manager/manager-web/src/main/resources/application.properties 
b/inlong-manager/manager-web/src/main/resources/application.properties
index 731a80d70..ed6f02848 100644
--- a/inlong-manager/manager-web/src/main/resources/application.properties
+++ b/inlong-manager/manager-web/src/main/resources/application.properties
@@ -25,7 +25,6 @@ server.servlet.context-path=/api/inlong/manager
 spring.application.name=InLong-Manager-Web
 spring.profiles.active=dev
 
-spring.main.allow-circular-references=true
 spring.mvc.pathmatch.matching-strategy=ANT_PATH_MATCHER
 
 # Serialize the Date type to a timestamp
diff --git 
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/WorkflowEngine.java
 
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/WorkflowEngine.java
index a1dc9cbc5..12aed2a97 100644
--- 
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/WorkflowEngine.java
+++ 
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/WorkflowEngine.java
@@ -29,13 +29,6 @@ public interface WorkflowEngine {
      */
     ProcessDefinitionService processDefinitionService();
 
-    /**
-     * Get process definition  repository
-     *
-     * @return Process repository.
-     */
-    ProcessDefinitionRepository processDefinitionRepository();
-
     /**
      * Get process instance service
      *
diff --git 
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/impl/WorkflowEngineImpl.java
 
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/impl/WorkflowEngineImpl.java
index c298ae84d..e4f6f3869 100644
--- 
a/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/impl/WorkflowEngineImpl.java
+++ 
b/inlong-manager/manager-workflow/src/main/java/org/apache/inlong/manager/workflow/core/impl/WorkflowEngineImpl.java
@@ -42,8 +42,6 @@ public class WorkflowEngineImpl implements WorkflowEngine {
 
     private final ProcessDefinitionService processDefService;
 
-    private final ProcessDefinitionRepository processDefRepository;
-
     private final ProcessService processService;
 
     private final TaskService taskService;
@@ -59,9 +57,6 @@ public class WorkflowEngineImpl implements WorkflowEngine {
         // Database transaction assistant
         TransactionHelper transactionHelper = new 
TransactionHelper(workflowConfig.getTransactionManager());
 
-        // Get workflow data accessor
-        this.processDefRepository = workflowConfig.getDefinitionRepository();
-
         // Workflow event listener manager
         EventListenerManagerFactory listenerManagerFactory = new 
EventListenerManagerFactory(workflowConfig);
 
@@ -74,6 +69,8 @@ public class WorkflowEngineImpl implements WorkflowEngine {
         ProcessorExecutor processorExecutor = new 
ProcessorExecutorImpl(processEntityMapper,
                 taskEntityMapper, workflowEventNotifier, transactionHelper);
 
+        ProcessDefinitionRepository processDefRepository = 
workflowConfig.getDefinitionRepository();
+
         // Workflow context builder
         WorkflowContextBuilder contextBuilder = new WorkflowContextBuilderImpl(
                 processDefRepository, processEntityMapper, taskEntityMapper);
@@ -100,11 +97,6 @@ public class WorkflowEngineImpl implements WorkflowEngine {
         return processDefService;
     }
 
-    @Override
-    public ProcessDefinitionRepository processDefinitionRepository() {
-        return processDefRepository;
-    }
-
     @Override
     public ProcessService processService() {
         return processService;

Reply via email to