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);