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 c58500501d [INLONG-8958][Manager] Fix the fail for creating cls topic 
(#9029)
c58500501d is described below

commit c58500501da4fc69d51616d60590d67818c72e29
Author: castor <[email protected]>
AuthorDate: Fri Oct 13 16:28:54 2023 +0800

    [INLONG-8958][Manager] Fix the fail for creating cls topic (#9029)
---
 .../inlong/manager/pojo/sink/cls/ClsSink.java      |   3 +
 .../inlong/manager/pojo/sink/cls/ClsSinkDTO.java   |   4 +-
 .../manager/pojo/sink/cls/ClsSinkRequest.java      |   4 +-
 .../service/core/impl/SortClusterServiceImpl.java  |   6 +-
 .../service/resource/sink/cls/ClsOperator.java     | 241 +++++++++++++++++++++
 .../resource/sink/cls/ClsResourceOperator.java     | 112 +++-------
 .../manager/service/sink/AbstractSinkOperator.java |  12 +-
 .../manager/service/sink/StreamSinkOperator.java   |   3 +-
 .../manager/service/sink/cls/ClsSinkOperator.java  |  18 +-
 .../service/sink/es/ElasticsearchSinkOperator.java |   5 +-
 .../service/sink/pulsar/PulsarSinkOperator.java    |   6 +-
 11 files changed, 308 insertions(+), 106 deletions(-)

diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSink.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSink.java
index 5ab92728fd..1668da82b8 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSink.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSink.java
@@ -79,6 +79,9 @@ public class ClsSink extends StreamSink {
     @ApiModelProperty("Cloud log service index tokenizer")
     private String tokenizer;
 
+    @ApiModelProperty("Cloud log service topic storage duration")
+    private Integer storageDuration;
+
     public ClsSink() {
         this.setSinkType(SinkType.CLS);
     }
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkDTO.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkDTO.java
index 3acd1e0372..f76f5bfb01 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkDTO.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkDTO.java
@@ -46,8 +46,8 @@ public class ClsSinkDTO {
     @ApiModelProperty("Cloud log service topic name")
     private String topicName;
 
-    @ApiModelProperty("Cloud log service topic save time")
-    private Integer saveTime;
+    @ApiModelProperty("Cloud log service topic storage duration")
+    private Integer storageDuration;
 
     @ApiModelProperty("Cloud log service tag name")
     private String tag;
diff --git 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkRequest.java
 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkRequest.java
index d962b3e9ac..7f109185c2 100644
--- 
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkRequest.java
+++ 
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sink/cls/ClsSinkRequest.java
@@ -40,8 +40,8 @@ public class ClsSinkRequest extends SinkRequest {
     @ApiModelProperty("Cloud log service topic name")
     private String topicName;
 
-    @ApiModelProperty("Cloud log service topic save time")
-    private Integer saveTime;
+    @ApiModelProperty("Cloud log service topic storage duration")
+    private Integer storageDuration;
 
     @ApiModelProperty("Cloud log service tag name")
     private String tag;
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortClusterServiceImpl.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortClusterServiceImpl.java
index ff7142d614..6af22b75b6 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortClusterServiceImpl.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortClusterServiceImpl.java
@@ -243,7 +243,7 @@ public class SortClusterServiceImpl implements 
SortClusterService {
                         return SortTaskConfig.builder()
                                 .name(taskName)
                                 .type(type)
-                                .idParams(this.parseIdParams(streams))
+                                .idParams(this.parseIdParams(streams, 
nodeInfo))
                                 .sinkParams(this.parseSinkParams(nodeInfo))
                                 .build();
                     } catch (Exception e) {
@@ -260,13 +260,13 @@ public class SortClusterServiceImpl implements 
SortClusterService {
                 .build();
     }
 
-    private List<Map<String, String>> parseIdParams(List<StreamSinkEntity> 
streams) {
+    private List<Map<String, String>> parseIdParams(List<StreamSinkEntity> 
streams, DataNodeInfo dataNodeInfo) {
         return streams.stream()
                 .map(streamSink -> {
                     try {
                         StreamSinkOperator operator = 
sinkOperatorFactory.getInstance(streamSink.getSinkType());
                         List<String> fields = 
fieldMap.get(streamSink.getInlongGroupId());
-                        return operator.parse2IdParams(streamSink, fields);
+                        return operator.parse2IdParams(streamSink, fields, 
dataNodeInfo);
                     } catch (Exception e) {
                         LOGGER.error("fail to parse id params of groupId={}, 
streamId={} name={}, type={}}",
                                 streamSink.getInlongGroupId(), 
streamSink.getInlongStreamId(),
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsOperator.java
new file mode 100644
index 0000000000..9aa29d9993
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsOperator.java
@@ -0,0 +1,241 @@
+/*
+ * 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.resource.sink.cls;
+
+import org.apache.inlong.manager.common.consts.InlongConstants;
+import org.apache.inlong.manager.common.exceptions.BusinessException;
+
+import com.tencentcloudapi.cls.v20201016.ClsClient;
+import com.tencentcloudapi.cls.v20201016.models.CreateIndexRequest;
+import com.tencentcloudapi.cls.v20201016.models.CreateIndexResponse;
+import com.tencentcloudapi.cls.v20201016.models.CreateTopicRequest;
+import com.tencentcloudapi.cls.v20201016.models.CreateTopicResponse;
+import com.tencentcloudapi.cls.v20201016.models.DescribeIndexRequest;
+import com.tencentcloudapi.cls.v20201016.models.DescribeIndexResponse;
+import com.tencentcloudapi.cls.v20201016.models.DescribeTopicsRequest;
+import com.tencentcloudapi.cls.v20201016.models.DescribeTopicsResponse;
+import com.tencentcloudapi.cls.v20201016.models.Filter;
+import com.tencentcloudapi.cls.v20201016.models.FullTextInfo;
+import com.tencentcloudapi.cls.v20201016.models.ModifyIndexRequest;
+import com.tencentcloudapi.cls.v20201016.models.ModifyIndexResponse;
+import com.tencentcloudapi.cls.v20201016.models.ModifyTopicRequest;
+import com.tencentcloudapi.cls.v20201016.models.ModifyTopicResponse;
+import com.tencentcloudapi.cls.v20201016.models.RuleInfo;
+import com.tencentcloudapi.cls.v20201016.models.Tag;
+import com.tencentcloudapi.cls.v20201016.models.TopicInfo;
+import com.tencentcloudapi.common.Credential;
+import com.tencentcloudapi.common.exception.TencentCloudSDKException;
+import com.tencentcloudapi.common.profile.ClientProfile;
+import com.tencentcloudapi.common.profile.HttpProfile;
+import org.apache.commons.lang3.ArrayUtils;
+import org.apache.commons.lang3.ObjectUtils;
+import org.apache.commons.lang3.StringUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.stereotype.Service;
+
+import java.util.ArrayList;
+import java.util.List;
+
+@Service
+public class ClsOperator {
+
+    private static final Logger LOG = 
LoggerFactory.getLogger(ClsOperator.class);
+    private static final String TOPIC_NAME = "topicName";
+    private static final String LOG_SET_ID = "logsetId";
+    private static final long PRECISE_SEARCH = 1L;
+
+    public String createTopicReturnTopicId(String topicName, String logSetId, 
String tag, String secretId,
+            String secretKey, String endPoint, String region)
+            throws TencentCloudSDKException {
+        ClsClient client = getClsClient(secretId, secretKey, endPoint, region);
+        CreateTopicRequest req = getCreateTopicRequest(tag, logSetId, 
topicName);
+        CreateTopicResponse resp = client.CreateTopic(req);
+        LOG.info("create cls topic success for topicName = {}, topicId = {}, 
requestId = {}", topicName,
+                resp.getTopicId(), resp.getRequestId());
+        updateTopicTag(resp.getTopicId(), tag, secretId, secretKey, endPoint, 
region);
+        return resp.getTopicId();
+    }
+
+    public void updateTopicTag(String topicId, String tag, String secretId,
+            String secretKey, String endPoint, String region) throws 
TencentCloudSDKException {
+        ClsClient client = getClsClient(secretId, secretKey, endPoint, region);
+        ModifyTopicRequest modifyTopicRequest = new ModifyTopicRequest();
+        
modifyTopicRequest.setTags(convertTags(tag.split(InlongConstants.CENTER_LINE)));
+        modifyTopicRequest.setTopicId(topicId);
+        ModifyTopicResponse resp = client.ModifyTopic(modifyTopicRequest);
+        LOG.info("update cls topic tag success for topicId = {}, requestId = 
{}", topicId, resp.getRequestId());
+    }
+
+    /**
+     * Create topic index by tokenizer
+     */
+    public void createTopicIndex(String tokenizer, String topicId, String 
secretId, String secretKey, String endPoint,
+            String region) throws BusinessException {
+
+        LOG.debug("create topic index start for topicId = {}, tokenizer = {}", 
topicId, tokenizer);
+        if (StringUtils.isBlank(tokenizer)) {
+            LOG.warn("tokenizer is blank for topic = {}", topicId);
+            return;
+        }
+        FullTextInfo topicIndexFullText = getTopicIndexFullText(secretId, 
secretKey, endPoint, region, topicId);
+        if (ObjectUtils.anyNotNull(topicIndexFullText)) {
+            // if topic index exist, update
+            LOG.debug("cls topic is exist and update for topicId = 
{},tokenizer = {}", topicId, tokenizer);
+            updateTopicIndex(tokenizer, topicId, secretId, secretKey, 
endPoint, region);
+            return;
+        }
+        ClsClient clsClient = getClsClient(secretId, secretKey, endPoint, 
region);
+        CreateIndexRequest req = getCreateIndexRequest(tokenizer, topicId);
+        try {
+            CreateIndexResponse createIndexResponse = 
clsClient.CreateIndex(req);
+            LOG.debug("create index success for topic = {}, tokenizer = {}, 
requestId = {}", topicId,
+                    tokenizer, createIndexResponse.getRequestId());
+        } catch (TencentCloudSDKException e) {
+            String errMsg = "Create cls topic index failed: " + e.getMessage();
+            LOG.error(errMsg, e);
+            throw new BusinessException(errMsg);
+        }
+
+    }
+
+    /**
+     * Describe cls topicId by topic name
+     */
+    public String describeTopicIDByTopicName(String topicName, String 
logSetId, String tag, String secretId,
+            String secretKey, String endPoint, String region) {
+        ClsClient clsClient = getClsClient(secretId, secretKey, endPoint, 
region);
+        Filter[] filters = getDescribeFilters(topicName, logSetId);
+        DescribeTopicsRequest req = new DescribeTopicsRequest();
+        req.setFilters(filters);
+        req.setPreciseSearch(PRECISE_SEARCH);
+        try {
+            DescribeTopicsResponse describeTopicsResponse = 
clsClient.DescribeTopics(req);
+            if (ArrayUtils.isNotEmpty(describeTopicsResponse.getTopics())) {
+                TopicInfo[] topics = describeTopicsResponse.getTopics();
+                return topics[0].getTopicId();
+            }
+            return null;
+        } catch (TencentCloudSDKException e) {
+            String errMsg = "describe cls topic failed: " + e.getMessage();
+            LOG.error(errMsg, e);
+            throw new BusinessException(errMsg);
+        }
+    }
+
+    public Filter[] getDescribeFilters(String topicName, String logSetId) {
+        Filter topicNameFilter = new Filter();
+        topicNameFilter.setKey(TOPIC_NAME);
+        String[] topicNameFilterValues = new String[]{topicName};
+        topicNameFilter.setValues(topicNameFilterValues);
+
+        Filter logSetIdFilter = new Filter();
+        logSetIdFilter.setKey(LOG_SET_ID);
+        String[] logSetFilterValues = new String[]{logSetId};
+        logSetIdFilter.setValues(logSetFilterValues);
+        return new Filter[]{topicNameFilter, logSetIdFilter};
+    }
+
+    /**
+     * Get cls topic index full text
+     */
+    public FullTextInfo getTopicIndexFullText(String secretId, String 
secretKey, String endPoint, String region,
+            String topicId) {
+
+        ClsClient clsClient = getClsClient(secretId, secretKey, endPoint, 
region);
+        DescribeIndexRequest req = new DescribeIndexRequest();
+        req.setTopicId(topicId);
+        try {
+            DescribeIndexResponse resp = clsClient.DescribeIndex(req);
+            return resp.getRule() == null ? null : 
resp.getRule().getFullText();
+        } catch (TencentCloudSDKException e) {
+            String errMsg = "describe cls topic index failed: " + 
e.getMessage();
+            LOG.error(errMsg, e);
+            throw new BusinessException(errMsg);
+        }
+    }
+
+    public void updateTopicIndex(String tokenizer, String topicId,
+            String secretId, String secretKey, String endPoint, String region) 
{
+        ClsClient clsClient = getClsClient(secretId, secretKey, endPoint, 
region);
+        RuleInfo ruleInfo = new RuleInfo();
+        FullTextInfo fullTextInfo = new FullTextInfo();
+        fullTextInfo.setTokenizer(tokenizer);
+        ruleInfo.setFullText(fullTextInfo);
+
+        ModifyIndexRequest req = new ModifyIndexRequest();
+        req.setTopicId(topicId);
+        req.setRule(ruleInfo);
+        try {
+            ModifyIndexResponse modifyIndexResponse = 
clsClient.ModifyIndex(req);
+            LOG.debug("update index success for topicId = {}, tokenizer = {}, 
requestId = {}", topicId, tokenizer,
+                    modifyIndexResponse.getRequestId());
+        } catch (TencentCloudSDKException e) {
+            String errMsg = "update cls topic index failed: " + e.getMessage();
+            LOG.error(errMsg, e);
+            throw new BusinessException(errMsg);
+        }
+    }
+
+    public ClsClient getClsClient(String secretId, String secretKey, String 
endPoint, String region) {
+        Credential cred = new Credential(secretId,
+                secretKey);
+        HttpProfile httpProfile = new HttpProfile();
+        httpProfile.setEndpoint(endPoint);
+        ClientProfile clientProfile = new ClientProfile();
+
+        clientProfile.setHttpProfile(httpProfile);
+        return new ClsClient(cred, region, clientProfile);
+    }
+
+    public CreateIndexRequest getCreateIndexRequest(String tokenizer, String 
topicId) {
+        RuleInfo ruleInfo = new RuleInfo();
+        FullTextInfo fullTextInfo = new FullTextInfo();
+        fullTextInfo.setTokenizer(tokenizer);
+        ruleInfo.setFullText(fullTextInfo);
+
+        CreateIndexRequest req = new CreateIndexRequest();
+        req.setTopicId(topicId);
+        req.setRule(ruleInfo);
+        return req;
+    }
+
+    public CreateTopicRequest getCreateTopicRequest(String tags, String 
logSetId, String topicName) {
+        CreateTopicRequest req = new CreateTopicRequest();
+        req.setTags(convertTags(tags.split(InlongConstants.CENTER_LINE)));
+        req.setLogsetId(logSetId);
+        req.setTopicName(topicName);
+        return req;
+    }
+
+    public Tag[] convertTags(String[] allTags) {
+        List<Tag> tagList = new ArrayList<>();
+        for (String tag : allTags) {
+            String[] keyAndValueOfTag = tag.split(InlongConstants.COLON);
+            if (keyAndValueOfTag.length < 2) {
+                continue;
+            }
+            Tag tagInfo = new Tag();
+            tagInfo.setKey(keyAndValueOfTag[0]);
+            tagInfo.setValue(keyAndValueOfTag[1]);
+            tagList.add(tagInfo);
+        }
+        return tagList.toArray(new Tag[0]);
+    }
+
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
index 20effc5f93..c16391907a 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
@@ -22,7 +22,6 @@ import 
org.apache.inlong.manager.common.consts.InlongConstants;
 import org.apache.inlong.manager.common.consts.SinkType;
 import org.apache.inlong.manager.common.enums.SinkStatus;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
-import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.JsonUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
@@ -34,17 +33,7 @@ import org.apache.inlong.manager.pojo.sink.cls.ClsSinkDTO;
 import 
org.apache.inlong.manager.service.resource.sink.AbstractStandaloneSinkResourceOperator;
 import org.apache.inlong.manager.service.sink.StreamSinkService;
 
-import com.tencentcloudapi.cls.v20201016.ClsClient;
-import com.tencentcloudapi.cls.v20201016.models.CreateIndexRequest;
-import com.tencentcloudapi.cls.v20201016.models.CreateTopicRequest;
-import com.tencentcloudapi.cls.v20201016.models.CreateTopicResponse;
-import com.tencentcloudapi.cls.v20201016.models.FullTextInfo;
-import com.tencentcloudapi.cls.v20201016.models.RuleInfo;
-import com.tencentcloudapi.cls.v20201016.models.Tag;
-import com.tencentcloudapi.common.Credential;
 import com.tencentcloudapi.common.exception.TencentCloudSDKException;
-import com.tencentcloudapi.common.profile.ClientProfile;
-import com.tencentcloudapi.common.profile.HttpProfile;
 import org.apache.commons.lang3.StringUtils;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -62,6 +51,8 @@ public class ClsResourceOperator extends 
AbstractStandaloneSinkResourceOperator
     private StreamSinkService sinkService;
     @Autowired
     private StreamSinkEntityMapper streamSinkEntityMapper;
+    @Autowired
+    private ClsOperator clsOperator;
 
     @Override
     public Boolean accept(String sinkType) {
@@ -78,101 +69,58 @@ public class ClsResourceOperator extends 
AbstractStandaloneSinkResourceOperator
             LOG.warn("create resource was disabled, skip to create for [" + 
sinkInfo.getId() + "]");
             return;
         }
-        this.createTopicID(sinkInfo);
+        this.createClsResource(sinkInfo);
         this.assignCluster(sinkInfo);
     }
 
     /**
      * Create cloud log service topic
      */
-    private void createTopicID(SinkInfo sinkInfo) {
+    private void createClsResource(SinkInfo sinkInfo) {
         ClsDataNodeDTO clsDataNode = getClsDataNode(sinkInfo);
         ClsSinkDTO clsSinkDTO = JsonUtils.parseObject(sinkInfo.getExtParams(), 
ClsSinkDTO.class);
         try {
-            ClsClient client = getClsClient(clsDataNode);
-            CreateTopicRequest req = getCreateTopicRequest(clsDataNode, 
clsSinkDTO);
-            CreateTopicResponse resp = client.CreateTopic(req);
-            LOG.info("create cls topic {} success ,topicId {}", 
clsSinkDTO.getTopicName(), resp.getTopicId());
-            // update set topic id into sink info
-            clsSinkDTO.setTopicId(resp.getTopicId());
+            String topicId = getTopicID(clsDataNode, clsSinkDTO);
+            clsSinkDTO.setTopicId(topicId);
             sinkInfo.setExtParams(JsonUtils.toJsonString(clsSinkDTO));
             // create topic index by tokenizer
-            this.createTopicIndex(sinkInfo);
-            StreamSinkEntity streamSinkEntity = new StreamSinkEntity();
-            CommonBeanUtils.copyProperties(sinkInfo, streamSinkEntity, true);
-            streamSinkEntityMapper.updateByIdSelective(streamSinkEntity);
+            clsOperator.createTopicIndex(clsSinkDTO.getTokenizer(), 
clsSinkDTO.getTopicId(),
+                    clsDataNode.getManageSecretId(),
+                    clsDataNode.getManageSecretKey(), 
clsDataNode.getEndpoint(), clsDataNode.getRegion());
+            // update set topic id into sink info
+            updateSinkInfo(sinkInfo, clsSinkDTO);
             String info = "success to create cls resource";
             sinkService.updateStatus(sinkInfo.getId(), 
SinkStatus.CONFIG_SUCCESSFUL.getCode(), info);
-            LOG.info("update cls sink = {}info status  success ,topicName {}", 
streamSinkEntity.getSinkName(),
+            LOG.info("update cls info status success for sinkId= {}, topicName 
= {}", sinkInfo.getSinkName(),
                     clsSinkDTO.getTopicName());
         } catch (TencentCloudSDKException e) {
-            String errMsg = "Create cls topic  failed: " + e.getMessage();
+            String errMsg = "Create cls topic failed: " + e.getMessage();
             LOG.error(errMsg, e);
             sinkService.updateStatus(sinkInfo.getId(), 
SinkStatus.CONFIG_FAILED.getCode(), errMsg);
             throw new BusinessException(errMsg);
         }
     }
 
-    private CreateTopicRequest getCreateTopicRequest(ClsDataNodeDTO 
clsDataNode, ClsSinkDTO clsSinkDTO) {
-        CreateTopicRequest req = new CreateTopicRequest();
-        String[] allTags = 
clsSinkDTO.getTag().split(InlongConstants.CENTER_LINE);
-        req.setTags(convertTags(allTags));
-        req.setLogsetId(clsDataNode.getLogSetId());
-        req.setTopicName(clsSinkDTO.getTopicName());
-        return req;
-    }
-
-    private ClsClient getClsClient(ClsDataNodeDTO clsDataNode) {
-        Credential cred = new Credential(clsDataNode.getManageSecretId(),
-                clsDataNode.getManageSecretId());
-        HttpProfile httpProfile = new HttpProfile();
-        httpProfile.setEndpoint(clsDataNode.getEndpoint());
-        ClientProfile clientProfile = new ClientProfile();
-
-        clientProfile.setHttpProfile(httpProfile);
-        return new ClsClient(cred, clsDataNode.getRegion(), clientProfile);
-    }
-
-    /**
-     * Create topic index by tokenizer
-     */
-    private void createTopicIndex(SinkInfo sinkInfo) throws BusinessException {
-        ClsSinkDTO clsSinkDTO = JsonUtils.parseObject(sinkInfo.getExtParams(), 
ClsSinkDTO.class);
-        if (StringUtils.isNotBlank(clsSinkDTO.getTokenizer())) {
-            LOG.warn("topic {} tokenizer is empty", clsSinkDTO.getTopicName());
-            return;
-        }
-        ClsDataNodeDTO clsDataNode = getClsDataNode(sinkInfo);
-        ClsClient clsClient = getClsClient(clsDataNode);
-        RuleInfo ruleInfo = new RuleInfo();
-        FullTextInfo fullTextInfo = new FullTextInfo();
-        fullTextInfo.setTokenizer(clsSinkDTO.getTokenizer());
-        ruleInfo.setFullText(fullTextInfo);
-
-        CreateIndexRequest req = new CreateIndexRequest();
-        req.setTopicId(clsSinkDTO.getTopicId());
-        req.setRule(ruleInfo);
-        try {
-            clsClient.CreateIndex(req);
-        } catch (TencentCloudSDKException e) {
-            String errMsg = "Create cls topic index failed: " + e.getMessage();
-            LOG.error(errMsg, e);
-            sinkService.updateStatus(sinkInfo.getId(), 
SinkStatus.CONFIG_FAILED.getCode(), errMsg);
-            throw new BusinessException(errMsg);
+    private String getTopicID(ClsDataNodeDTO clsDataNode, ClsSinkDTO 
clsSinkDTO)
+            throws TencentCloudSDKException {
+        String topicId = 
clsOperator.describeTopicIDByTopicName(clsSinkDTO.getTopicName(), 
clsDataNode.getLogSetId(),
+                clsSinkDTO.getTag(),
+                clsDataNode.getManageSecretId(), 
clsDataNode.getManageSecretKey(), clsDataNode.getEndpoint(),
+                clsDataNode.getRegion());
+        if (StringUtils.isBlank(topicId)) {
+            // if topic don't exist, create topic in cls
+            topicId = 
clsOperator.createTopicReturnTopicId(clsSinkDTO.getTopicName(), 
clsDataNode.getLogSetId(),
+                    clsSinkDTO.getTag(), clsDataNode.getManageSecretId(), 
clsDataNode.getManageSecretKey(),
+                    clsDataNode.getEndpoint(),
+                    clsDataNode.getRegion());
         }
-        LOG.info("topic {} create index success tokenizer is {}", 
clsSinkDTO.getTopicName(), clsSinkDTO.getTokenizer());
+        return topicId;
     }
 
-    private Tag[] convertTags(String[] allTags) {
-        Tag[] tags = new Tag[allTags.length];
-        for (int i = 0; i < allTags.length; i++) {
-            String tag = allTags[i];
-            String[] keyAndValueOfTag = tag.split(InlongConstants.COLON);
-            Tag tagInfo = new Tag();
-            tagInfo.set(keyAndValueOfTag[0], keyAndValueOfTag[1]);
-            tags[i] = tagInfo;
-        }
-        return tags;
+    private void updateSinkInfo(SinkInfo sinkInfo, ClsSinkDTO clsSinkDTO) {
+        StreamSinkEntity streamSinkEntity = 
streamSinkEntityMapper.selectByPrimaryKey(sinkInfo.getId());
+        streamSinkEntity.setExtParams(JsonUtils.toJsonString(clsSinkDTO));
+        streamSinkEntityMapper.updateByIdSelective(streamSinkEntity);
     }
 
     private ClsDataNodeDTO getClsDataNode(SinkInfo sinkInfo) {
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
index 57b241cfcf..59a5dbe432 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/AbstractSinkOperator.java
@@ -29,6 +29,7 @@ import 
org.apache.inlong.manager.dao.mapper.InlongStreamEntityMapper;
 import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
 import org.apache.inlong.manager.dao.mapper.StreamSinkFieldEntityMapper;
 import org.apache.inlong.manager.pojo.common.PageResult;
+import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.sink.SinkField;
 import org.apache.inlong.manager.pojo.sink.SinkRequest;
 import org.apache.inlong.manager.pojo.sink.StreamSink;
@@ -218,12 +219,17 @@ public abstract class AbstractSinkOperator implements 
StreamSinkOperator {
     }
 
     @Override
-    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields) {
+    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields,
+            DataNodeInfo dataNodeInfo) {
         Map<String, String> param;
         try {
-            param = JsonUtils.parseObject(streamSink.getExtParams(), 
HashMap.class);
+            HashMap<String, Object> streamInfoMap = 
JsonUtils.parseObject(streamSink.getExtParams(), HashMap.class);
+            param = new HashMap<>();
+            assert streamInfoMap != null;
+            for (String key : streamInfoMap.keySet()) {
+                param.put(key, String.valueOf(streamInfoMap.get(key)));
+            }
             // put group and stream info
-            assert param != null;
             param.put(KEY_GROUP_ID, streamSink.getInlongGroupId());
             param.put(KEY_STREAM_ID, streamSink.getInlongStreamId());
             return param;
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
index fe5e14987a..619d7abed8 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/StreamSinkOperator.java
@@ -20,6 +20,7 @@ package org.apache.inlong.manager.service.sink;
 import org.apache.inlong.manager.common.enums.SinkStatus;
 import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
 import org.apache.inlong.manager.pojo.common.PageResult;
+import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.sink.SinkField;
 import org.apache.inlong.manager.pojo.sink.SinkRequest;
 import org.apache.inlong.manager.pojo.sink.StreamSink;
@@ -115,5 +116,5 @@ public interface StreamSinkOperator {
      * @param streamSink
      * @return
      */
-    Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields);
+    Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields, DataNodeInfo dataNodeInfo);
 }
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/cls/ClsSinkOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/cls/ClsSinkOperator.java
index ad8f4c69c1..a2ec9e1958 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/cls/ClsSinkOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/cls/ClsSinkOperator.java
@@ -27,7 +27,9 @@ import org.apache.inlong.manager.common.util.JsonUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
 import org.apache.inlong.manager.dao.mapper.DataNodeEntityMapper;
+import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.cls.ClsDataNodeDTO;
+import org.apache.inlong.manager.pojo.node.cls.ClsDataNodeInfo;
 import org.apache.inlong.manager.pojo.sink.SinkField;
 import org.apache.inlong.manager.pojo.sink.SinkRequest;
 import org.apache.inlong.manager.pojo.sink.StreamSink;
@@ -110,17 +112,15 @@ public class ClsSinkOperator extends AbstractSinkOperator 
{
     }
 
     @Override
-    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields) {
-        Map<String, String> params = super.parse2IdParams(streamSink, fields);
+    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields,
+            DataNodeInfo dataNodeInfo) {
+        Map<String, String> params = super.parse2IdParams(streamSink, fields, 
dataNodeInfo);
         ClsSinkDTO clsSinkDTO = 
JsonUtils.parseObject(streamSink.getExtParams(), ClsSinkDTO.class);
         params.put(TOPIC_ID, clsSinkDTO.getTopicId());
-        DataNodeEntity dataNodeEntity = 
dataNodeEntityMapper.selectByUniqueKey(streamSink.getDataNodeName(),
-                DataNodeType.CLS);
-        ClsDataNodeDTO clsDataNodeDTO = 
JsonUtils.parseObject(dataNodeEntity.getExtParams(),
-                ClsDataNodeDTO.class);
-        params.put(SECRET_ID, clsDataNodeDTO.getSendSecretId());
-        params.put(SECRET_KEY, clsDataNodeDTO.getSendSecretKey());
-        params.put(END_POINT, clsDataNodeDTO.getEndpoint());
+        ClsDataNodeInfo clsDataNodeInfo = (ClsDataNodeInfo) dataNodeInfo;
+        params.put(SECRET_ID, clsDataNodeInfo.getSendSecretId());
+        params.put(SECRET_KEY, clsDataNodeInfo.getSendSecretKey());
+        params.put(END_POINT, clsDataNodeInfo.getEndpoint());
         StringBuilder fieldNames = new StringBuilder();
         for (String field : fields) {
             fieldNames.append(field).append(InlongConstants.BLANK);
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/es/ElasticsearchSinkOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/es/ElasticsearchSinkOperator.java
index 4946e34d53..06ea6fac5f 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/es/ElasticsearchSinkOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/es/ElasticsearchSinkOperator.java
@@ -118,8 +118,9 @@ public class ElasticsearchSinkOperator extends 
AbstractSinkOperator {
     }
 
     @Override
-    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields) {
-        Map<String, String> idParams = super.parse2IdParams(streamSink, 
fields);
+    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields,
+            DataNodeInfo dataNodeInfo) {
+        Map<String, String> idParams = super.parse2IdParams(streamSink, 
fields, dataNodeInfo);
         StringBuilder sb = new StringBuilder();
         for (String field : fields) {
             sb.append(field).append(" ");
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/pulsar/PulsarSinkOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/pulsar/PulsarSinkOperator.java
index 9452d6a73e..6a1657b1ec 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/pulsar/PulsarSinkOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sink/pulsar/PulsarSinkOperator.java
@@ -26,6 +26,7 @@ import org.apache.inlong.manager.common.util.JsonUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
 import org.apache.inlong.manager.dao.mapper.DataNodeEntityMapper;
+import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.pulsar.PulsarDataNodeDTO;
 import org.apache.inlong.manager.pojo.sink.SinkField;
 import org.apache.inlong.manager.pojo.sink.SinkRequest;
@@ -105,9 +106,10 @@ public class PulsarSinkOperator extends 
AbstractSinkOperator {
     }
 
     @Override
-    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields) {
+    public Map<String, String> parse2IdParams(StreamSinkEntity streamSink, 
List<String> fields,
+            DataNodeInfo dataNodeInfo) {
 
-        Map<String, String> params = super.parse2IdParams(streamSink, fields);
+        Map<String, String> params = super.parse2IdParams(streamSink, fields, 
dataNodeInfo);
         PulsarSinkDTO pulsarSinkDTO;
         try {
             pulsarSinkDTO = objectMapper.readValue(streamSink.getExtParams(), 
PulsarSinkDTO.class);


Reply via email to