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;