This is an automated email from the ASF dual-hosted git repository.
dockerzhang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git
The following commit(s) were added to refs/heads/master by this push:
new 9bedbac80 [INLONG-7706][Manager] Fix sink is always in the
configuration after being saved (#7707)
9bedbac80 is described below
commit 9bedbac80daff2894e6ed9bcb1ed63e95713a369
Author: fuweng11 <[email protected]>
AuthorDate: Mon Mar 27 18:42:59 2023 +0800
[INLONG-7706][Manager] Fix sink is always in the configuration after being
saved (#7707)
---
.../service/sink/StreamSinkServiceImpl.java | 30 ++++++++++++++--------
1 file changed, 20 insertions(+), 10 deletions(-)
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkServiceImpl.java
index dcd8d1d1b..6ed96ee84 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkServiceImpl.java
@@ -137,7 +137,10 @@ public class StreamSinkServiceImpl implements
StreamSinkService {
if (streamSuccess ||
StreamStatus.CONFIG_FAILED.getCode().equals(streamEntity.getStatus())) {
boolean enableCreateResource =
InlongConstants.ENABLE_CREATE_RESOURCE.equals(
request.getEnableCreateResource());
- SinkStatus nextStatus = enableCreateResource ?
SinkStatus.CONFIG_ING : SinkStatus.CONFIG_SUCCESSFUL;
+ SinkStatus nextStatus = request.getStartProcess() ?
SinkStatus.CONFIG_ING : SinkStatus.NEW;
+ if (!enableCreateResource) {
+ nextStatus = SinkStatus.CONFIG_SUCCESSFUL;
+ }
StreamSinkEntity sinkEntity = sinkMapper.selectByPrimaryKey(id);
sinkEntity.setStatus(nextStatus.getCode());
sinkMapper.updateStatus(sinkEntity);
@@ -196,7 +199,10 @@ public class StreamSinkServiceImpl implements
StreamSinkService {
if (streamSuccess ||
StreamStatus.CONFIG_FAILED.getCode().equals(streamEntity.getStatus())) {
boolean enableCreateResource =
InlongConstants.ENABLE_CREATE_RESOURCE.equals(
request.getEnableCreateResource());
- SinkStatus nextStatus = enableCreateResource ?
SinkStatus.CONFIG_ING : SinkStatus.CONFIG_SUCCESSFUL;
+ SinkStatus nextStatus = request.getStartProcess() ?
SinkStatus.CONFIG_ING : SinkStatus.NEW;
+ if (!enableCreateResource) {
+ nextStatus = SinkStatus.CONFIG_SUCCESSFUL;
+ }
StreamSinkEntity sinkEntity = sinkMapper.selectByPrimaryKey(id);
sinkEntity.setStatus(nextStatus.getCode());
sinkMapper.updateStatus(sinkEntity);
@@ -379,8 +385,9 @@ public class StreamSinkServiceImpl implements
StreamSinkService {
}
SinkStatus nextStatus = null;
- boolean streamSuccess =
StreamStatus.CONFIG_SUCCESSFUL.getCode().equals(streamEntity.getStatus());
- if (streamSuccess ||
StreamStatus.CONFIG_FAILED.getCode().equals(streamEntity.getStatus())) {
+ boolean enableConfig =
StreamStatus.CONFIG_SUCCESSFUL.getCode().equals(streamEntity.getStatus())
+ ||
StreamStatus.CONFIG_FAILED.getCode().equals(streamEntity.getStatus());
+ if (enableConfig) {
boolean enableCreateResource =
InlongConstants.ENABLE_CREATE_RESOURCE.equals(
request.getEnableCreateResource());
nextStatus = enableCreateResource ? SinkStatus.CONFIG_ING :
SinkStatus.CONFIG_SUCCESSFUL;
@@ -388,8 +395,9 @@ public class StreamSinkServiceImpl implements
StreamSinkService {
StreamSinkOperator sinkOperator =
operatorFactory.getInstance(request.getSinkType());
sinkOperator.updateOpt(request, nextStatus, operator);
- // If the stream is [CONFIG_SUCCESSFUL], then asynchronously start the
[CREATE_STREAM_RESOURCE] process
- if (streamSuccess && request.getStartProcess()) {
+ // If the stream is [CONFIG_SUCCESSFUL] or [CONFIG_FAILED], then
asynchronously start the
+ // [CREATE_STREAM_RESOURCE] process
+ if (enableConfig && request.getStartProcess()) {
this.startProcessForSink(request.getInlongGroupId(),
request.getInlongStreamId(), operator);
}
@@ -440,16 +448,18 @@ public class StreamSinkServiceImpl implements
StreamSinkService {
}
// update record
SinkStatus nextStatus = null;
- boolean streamSuccess =
StreamStatus.CONFIG_SUCCESSFUL.getCode().equals(streamEntity.getStatus());
- if (streamSuccess ||
StreamStatus.CONFIG_FAILED.getCode().equals(streamEntity.getStatus())) {
+ boolean enableConfig =
StreamStatus.CONFIG_SUCCESSFUL.getCode().equals(streamEntity.getStatus())
+ ||
StreamStatus.CONFIG_FAILED.getCode().equals(streamEntity.getStatus());
+ if (enableConfig) {
boolean enableCreateResource =
InlongConstants.ENABLE_CREATE_RESOURCE.equals(
request.getEnableCreateResource());
nextStatus = enableCreateResource ? SinkStatus.CONFIG_ING :
SinkStatus.CONFIG_SUCCESSFUL;
}
StreamSinkOperator sinkOperator =
operatorFactory.getInstance(request.getSinkType());
sinkOperator.updateOpt(request, nextStatus, opInfo.getName());
- // If the stream is [CONFIG_SUCCESSFUL], then asynchronously start the
[CREATE_STREAM_RESOURCE] process
- if (streamSuccess && request.getStartProcess()) {
+ // If the stream is [CONFIG_SUCCESSFUL] or [CONFIG_FAILED], then
asynchronously start the
+ // [CREATE_STREAM_RESOURCE] process
+ if (enableConfig && request.getStartProcess()) {
this.startProcessForSink(request.getInlongGroupId(),
request.getInlongStreamId(), opInfo.getName());
}
return true;