From 2f1ca2ad38b3c70e4809f32c078037f70ee5d6aa Mon Sep 17 00:00:00 2001 From: castorqin Date: Tue, 17 Oct 2023 20:43:55 +0800 Subject: [PATCH 01/18] [Manager] update StreamSinkFieldEntityMapper.selectAllFields method,use table inlong_stream_field replace stream_sink_field table --- .../mappers/StreamSinkFieldEntityMapper.xml | 2 +- .../core/impl/SortClusterServiceImpl.java | 55 ++++++++++++------- 2 files changed, 37 insertions(+), 20 deletions(-) diff --git a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml index 8273c333437..bc4a2bc127b 100644 --- a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml +++ b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml @@ -124,7 +124,7 @@ inlong_group_id, inlong_stream_id, field_name - from stream_sink_field + from inlong_stream_field where is_deleted = 0 order by id asc 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 6af22b75b6f..456ef46592a 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 @@ -17,24 +17,38 @@ package org.apache.inlong.manager.service.core.impl; +import com.google.gson.Gson; +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; +import javax.annotation.PostConstruct; +import org.apache.commons.codec.digest.DigestUtils; +import org.apache.commons.lang3.StringUtils; import org.apache.inlong.common.pojo.sortstandalone.SortClusterConfig; import org.apache.inlong.common.pojo.sortstandalone.SortClusterResponse; import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig; +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.pojo.node.DataNodeInfo; import org.apache.inlong.manager.pojo.sort.standalone.SortFieldInfo; +import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo; import org.apache.inlong.manager.pojo.sort.standalone.SortTaskInfo; +import org.apache.inlong.manager.pojo.stream.InlongStreamExtParam; import org.apache.inlong.manager.service.core.SortClusterService; import org.apache.inlong.manager.service.core.SortConfigLoader; import org.apache.inlong.manager.service.node.DataNodeOperator; import org.apache.inlong.manager.service.node.DataNodeOperatorFactory; import org.apache.inlong.manager.service.sink.SinkOperatorFactory; import org.apache.inlong.manager.service.sink.StreamSinkOperator; - -import com.google.gson.Gson; -import org.apache.commons.codec.digest.DigestUtils; -import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; @@ -42,20 +56,6 @@ import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; -import javax.annotation.PostConstruct; - -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Objects; -import java.util.Optional; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; - /** * Used to cache the sort cluster config and reduce the number of query to database. */ @@ -84,6 +84,8 @@ public class SortClusterServiceImpl implements SortClusterService { private Map sortClusterConfigMap = new ConcurrentHashMap<>(); // key : sort cluster name, value : error log private Map sortClusterErrorLogMap = new ConcurrentHashMap<>(); + // key: group id ,value: {key: stream id, value: stream info} + private Map> allStreams; private long reloadInterval; @@ -183,6 +185,12 @@ private void reloadAllClusterConfig() { && StringUtils.isNotBlank(dto.getSinkType())) .collect(Collectors.groupingBy(SortTaskInfo::getSortClusterName)); + // reload all streams + allStreams = sortConfigLoader.loadAllStreams() + .stream() + .collect(Collectors.groupingBy(SortSourceStreamInfo::getInlongGroupId, + Collectors.toMap(SortSourceStreamInfo::getInlongStreamId, info -> info))); + // get all stream sinks Map> task2AllStreams = sinkEntities.stream() .filter(entity -> StringUtils.isNotBlank(entity.getInlongClusterName())) @@ -266,7 +274,16 @@ private List> parseIdParams(List streams, try { StreamSinkOperator operator = sinkOperatorFactory.getInstance(streamSink.getSinkType()); List fields = fieldMap.get(streamSink.getInlongGroupId()); - return operator.parse2IdParams(streamSink, fields, dataNodeInfo); + Map params = operator.parse2IdParams(streamSink, fields, dataNodeInfo); + SortSourceStreamInfo sortSourceStreamInfo = allStreams.get(streamSink.getInlongGroupId()) + .get(streamSink.getInlongStreamId()); + InlongStreamExtParam inlongStreamExtParam = JsonUtils.parseObject( + sortSourceStreamInfo.getExtParams(), InlongStreamExtParam.class); + assert inlongStreamExtParam != null; + if(!inlongStreamExtParam.getUseExtendedFields()){ + params.put("fieldOffset", String.valueOf(0)); + } + return params; } catch (Exception e) { LOGGER.error("fail to parse id params of groupId={}, streamId={} name={}, type={}}", streamSink.getInlongGroupId(), streamSink.getInlongStreamId(), From 01fba8a67e29901ac8261122f482df1857bb3f16 Mon Sep 17 00:00:00 2001 From: castorqin Date: Wed, 18 Oct 2023 11:10:55 +0800 Subject: [PATCH 02/18] [Manager] update StreamSinkFieldEntityMapper.selectAllFields method,use table inlong_stream_field replace stream_sink_field table --- .../core/impl/SortClusterServiceImpl.java | 56 +++++++++++-------- 1 file changed, 33 insertions(+), 23 deletions(-) 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 456ef46592a..77e541d471d 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 @@ -17,21 +17,6 @@ package org.apache.inlong.manager.service.core.impl; -import com.google.gson.Gson; -import java.util.ArrayList; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Objects; -import java.util.Optional; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.Executors; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; -import java.util.stream.Collectors; -import javax.annotation.PostConstruct; -import org.apache.commons.codec.digest.DigestUtils; -import org.apache.commons.lang3.StringUtils; import org.apache.inlong.common.pojo.sortstandalone.SortClusterConfig; import org.apache.inlong.common.pojo.sortstandalone.SortClusterResponse; import org.apache.inlong.common.pojo.sortstandalone.SortTaskConfig; @@ -49,6 +34,10 @@ import org.apache.inlong.manager.service.node.DataNodeOperatorFactory; import org.apache.inlong.manager.service.sink.SinkOperatorFactory; import org.apache.inlong.manager.service.sink.StreamSinkOperator; + +import com.google.gson.Gson; +import org.apache.commons.codec.digest.DigestUtils; +import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.beans.factory.annotation.Autowired; @@ -56,6 +45,20 @@ import org.springframework.stereotype.Service; import org.springframework.transaction.annotation.Transactional; +import javax.annotation.PostConstruct; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; + /** * Used to cache the sort cluster config and reduce the number of query to database. */ @@ -76,6 +79,7 @@ public class SortClusterServiceImpl implements SortClusterService { private static final String KEY_GROUP_ID = "inlongGroupId"; private static final String KEY_STREAM_ID = "inlongStreamId"; + private static final String FILED_OFFSET = "fieldOffset"; private Map> fieldMap; // key : sort cluster name, value : md5 @@ -275,14 +279,7 @@ private List> parseIdParams(List streams, StreamSinkOperator operator = sinkOperatorFactory.getInstance(streamSink.getSinkType()); List fields = fieldMap.get(streamSink.getInlongGroupId()); Map params = operator.parse2IdParams(streamSink, fields, dataNodeInfo); - SortSourceStreamInfo sortSourceStreamInfo = allStreams.get(streamSink.getInlongGroupId()) - .get(streamSink.getInlongStreamId()); - InlongStreamExtParam inlongStreamExtParam = JsonUtils.parseObject( - sortSourceStreamInfo.getExtParams(), InlongStreamExtParam.class); - assert inlongStreamExtParam != null; - if(!inlongStreamExtParam.getUseExtendedFields()){ - params.put("fieldOffset", String.valueOf(0)); - } + params = setFiledOffset(streamSink, params); return params; } catch (Exception e) { LOGGER.error("fail to parse id params of groupId={}, streamId={} name={}, type={}}", @@ -295,6 +292,19 @@ private List> parseIdParams(List streams, .collect(Collectors.toList()); } + private Map setFiledOffset(StreamSinkEntity streamSink, Map params) { + + SortSourceStreamInfo sortSourceStreamInfo = allStreams.get(streamSink.getInlongGroupId()) + .get(streamSink.getInlongStreamId()); + InlongStreamExtParam inlongStreamExtParam = JsonUtils.parseObject( + sortSourceStreamInfo.getExtParams(), InlongStreamExtParam.class); + assert inlongStreamExtParam != null; + if (!inlongStreamExtParam.getUseExtendedFields()) { + params.put(FILED_OFFSET, String.valueOf(0)); + } + return params; + } + private Map parseSinkParams(DataNodeInfo nodeInfo) { DataNodeOperator operator = dataNodeOperatorFactory.getInstance(nodeInfo.getType()); return operator.parse2SinkParams(nodeInfo); From 1df5f262d092b700fbdb893ea8af013d72a882b0 Mon Sep 17 00:00:00 2001 From: castorqin Date: Wed, 18 Oct 2023 11:35:39 +0800 Subject: [PATCH 03/18] [Manager] update StreamSinkFieldEntityMapper.selectAllFields method,use table inlong_stream_field replace stream_sink_field table --- .../manager/service/core/impl/SortClusterServiceImpl.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) 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 77e541d471d..3c8dcba0e0f 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 @@ -298,8 +298,7 @@ private Map setFiledOffset(StreamSinkEntity streamSink, Map Date: Wed, 18 Oct 2023 15:12:22 +0800 Subject: [PATCH 04/18] [Manager] add default cls manager endpoint,use default cls endpoint to operator cls resource --- .../resource/sink/cls/ClsOperator.java | 39 ++++++++++--------- .../sink/cls/ClsResourceOperator.java | 6 +-- 2 files changed, 23 insertions(+), 22 deletions(-) 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 index 9aa29d99936..e366fcae78d 100644 --- 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 @@ -47,6 +47,7 @@ import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.springframework.beans.factory.annotation.Value; import org.springframework.stereotype.Service; import java.util.ArrayList; @@ -55,26 +56,28 @@ @Service public class ClsOperator { + @Value("${cls.manager.endpoint:cls.internal.tencentcloudapi.com}") + private String endpoint; 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) + String secretKey, String region) throws TencentCloudSDKException { - ClsClient client = getClsClient(secretId, secretKey, endPoint, region); + ClsClient client = getClsClient(secretId, secretKey, 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); + updateTopicTag(resp.getTopicId(), tag, secretId, secretKey, 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); + String secretKey, String region) throws TencentCloudSDKException { + ClsClient client = getClsClient(secretId, secretKey, region); ModifyTopicRequest modifyTopicRequest = new ModifyTopicRequest(); modifyTopicRequest.setTags(convertTags(tag.split(InlongConstants.CENTER_LINE))); modifyTopicRequest.setTopicId(topicId); @@ -85,7 +88,7 @@ public void updateTopicTag(String topicId, String tag, String secretId, /** * Create topic index by tokenizer */ - public void createTopicIndex(String tokenizer, String topicId, String secretId, String secretKey, String endPoint, + public void createTopicIndex(String tokenizer, String topicId, String secretId, String secretKey, String region) throws BusinessException { LOG.debug("create topic index start for topicId = {}, tokenizer = {}", topicId, tokenizer); @@ -93,14 +96,14 @@ public void createTopicIndex(String tokenizer, String topicId, String secretId, LOG.warn("tokenizer is blank for topic = {}", topicId); return; } - FullTextInfo topicIndexFullText = getTopicIndexFullText(secretId, secretKey, endPoint, region, topicId); + FullTextInfo topicIndexFullText = getTopicIndexFullText(secretId, secretKey, 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); + updateTopicIndex(tokenizer, topicId, secretId, secretKey, region); return; } - ClsClient clsClient = getClsClient(secretId, secretKey, endPoint, region); + ClsClient clsClient = getClsClient(secretId, secretKey, region); CreateIndexRequest req = getCreateIndexRequest(tokenizer, topicId); try { CreateIndexResponse createIndexResponse = clsClient.CreateIndex(req); @@ -117,9 +120,9 @@ public void createTopicIndex(String tokenizer, String topicId, String secretId, /** * 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); + public String describeTopicIDByTopicName(String topicName, String logSetId, String secretId, + String secretKey, String region) { + ClsClient clsClient = getClsClient(secretId, secretKey, region); Filter[] filters = getDescribeFilters(topicName, logSetId); DescribeTopicsRequest req = new DescribeTopicsRequest(); req.setFilters(filters); @@ -154,10 +157,10 @@ public Filter[] getDescribeFilters(String topicName, String logSetId) { /** * Get cls topic index full text */ - public FullTextInfo getTopicIndexFullText(String secretId, String secretKey, String endPoint, String region, + public FullTextInfo getTopicIndexFullText(String secretId, String secretKey, String region, String topicId) { - ClsClient clsClient = getClsClient(secretId, secretKey, endPoint, region); + ClsClient clsClient = getClsClient(secretId, secretKey, region); DescribeIndexRequest req = new DescribeIndexRequest(); req.setTopicId(topicId); try { @@ -171,8 +174,8 @@ public FullTextInfo getTopicIndexFullText(String secretId, String secretKey, Str } public void updateTopicIndex(String tokenizer, String topicId, - String secretId, String secretKey, String endPoint, String region) { - ClsClient clsClient = getClsClient(secretId, secretKey, endPoint, region); + String secretId, String secretKey, String region) { + ClsClient clsClient = getClsClient(secretId, secretKey, region); RuleInfo ruleInfo = new RuleInfo(); FullTextInfo fullTextInfo = new FullTextInfo(); fullTextInfo.setTokenizer(tokenizer); @@ -192,11 +195,11 @@ public void updateTopicIndex(String tokenizer, String topicId, } } - public ClsClient getClsClient(String secretId, String secretKey, String endPoint, String region) { + public ClsClient getClsClient(String secretId, String secretKey, String region) { Credential cred = new Credential(secretId, secretKey); HttpProfile httpProfile = new HttpProfile(); - httpProfile.setEndpoint(endPoint); + httpProfile.setEndpoint(endpoint); ClientProfile clientProfile = new ClientProfile(); clientProfile.setHttpProfile(httpProfile); 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 c16391907a4..b9e5c044a9e 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 @@ -86,7 +86,7 @@ private void createClsResource(SinkInfo sinkInfo) { // create topic index by tokenizer clsOperator.createTopicIndex(clsSinkDTO.getTokenizer(), clsSinkDTO.getTopicId(), clsDataNode.getManageSecretId(), - clsDataNode.getManageSecretKey(), clsDataNode.getEndpoint(), clsDataNode.getRegion()); + clsDataNode.getManageSecretKey(), clsDataNode.getRegion()); // update set topic id into sink info updateSinkInfo(sinkInfo, clsSinkDTO); String info = "success to create cls resource"; @@ -104,14 +104,12 @@ private void createClsResource(SinkInfo sinkInfo) { 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.getManageSecretId(), clsDataNode.getManageSecretKey(), 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()); } return topicId; From 22973d70c907ab52f753679d2ec8b50af0966204 Mon Sep 17 00:00:00 2001 From: castorqin Date: Wed, 18 Oct 2023 16:06:34 +0800 Subject: [PATCH 05/18] [Manager] add default cls manager endpoint,use default cls endpoint to operator cls resource --- .../inlong/manager/service/resource/sink/cls/ClsOperator.java | 2 +- .../manager/service/resource/sink/cls/ClsResourceOperator.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) 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 index e366fcae78d..81912b91ccb 100644 --- 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 @@ -63,7 +63,7 @@ public class ClsOperator { 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, + public String createTopicReturnTopicId(String topicName, String logSetId, String tag,Integer storageDuration, String secretId, String secretKey, String region) throws TencentCloudSDKException { ClsClient client = getClsClient(secretId, secretKey, region); 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 b9e5c044a9e..be1b5d1e836 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 @@ -109,7 +109,7 @@ private String getTopicID(ClsDataNodeDTO clsDataNode, ClsSinkDTO clsSinkDTO) 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(), + clsSinkDTO.getTag(),clsSinkDTO.getStorageDuration(), clsDataNode.getManageSecretId(), clsDataNode.getManageSecretKey(), clsDataNode.getRegion()); } return topicId; From f5655d14c95d2e978788d092f13508a04f05aaf7 Mon Sep 17 00:00:00 2001 From: castorqin Date: Wed, 18 Oct 2023 16:18:50 +0800 Subject: [PATCH 06/18] [Manager] add default cls manager endpoint,use default cls endpoint to operator cls resource --- .../manager/service/resource/sink/cls/ClsOperator.java | 9 ++++++--- .../service/resource/sink/cls/ClsResourceOperator.java | 3 ++- 2 files changed, 8 insertions(+), 4 deletions(-) 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 index 81912b91ccb..d0c85a6a53a 100644 --- 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 @@ -63,11 +63,12 @@ public class ClsOperator { private static final String LOG_SET_ID = "logsetId"; private static final long PRECISE_SEARCH = 1L; - public String createTopicReturnTopicId(String topicName, String logSetId, String tag,Integer storageDuration, String secretId, + public String createTopicReturnTopicId(String topicName, String logSetId, String tag, Integer storageDuration, + String secretId, String secretKey, String region) throws TencentCloudSDKException { ClsClient client = getClsClient(secretId, secretKey, region); - CreateTopicRequest req = getCreateTopicRequest(tag, logSetId, topicName); + CreateTopicRequest req = getCreateTopicRequest(tag, logSetId, topicName, storageDuration); CreateTopicResponse resp = client.CreateTopic(req); LOG.info("create cls topic success for topicName = {}, topicId = {}, requestId = {}", topicName, resp.getTopicId(), resp.getRequestId()); @@ -218,11 +219,13 @@ public CreateIndexRequest getCreateIndexRequest(String tokenizer, String topicId return req; } - public CreateTopicRequest getCreateTopicRequest(String tags, String logSetId, String topicName) { + public CreateTopicRequest getCreateTopicRequest(String tags, String logSetId, String topicName, + Integer storageDuration) { CreateTopicRequest req = new CreateTopicRequest(); req.setTags(convertTags(tags.split(InlongConstants.CENTER_LINE))); req.setLogsetId(logSetId); req.setTopicName(topicName); + req.setPeriod(storageDuration == null ? null : Long.valueOf(storageDuration)); return req; } 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 be1b5d1e836..173a139758f 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 @@ -109,7 +109,8 @@ private String getTopicID(ClsDataNodeDTO clsDataNode, ClsSinkDTO clsSinkDTO) if (StringUtils.isBlank(topicId)) { // if topic don't exist, create topic in cls topicId = clsOperator.createTopicReturnTopicId(clsSinkDTO.getTopicName(), clsDataNode.getLogSetId(), - clsSinkDTO.getTag(),clsSinkDTO.getStorageDuration(), clsDataNode.getManageSecretId(), clsDataNode.getManageSecretKey(), + clsSinkDTO.getTag(), clsSinkDTO.getStorageDuration(), clsDataNode.getManageSecretId(), + clsDataNode.getManageSecretKey(), clsDataNode.getRegion()); } return topicId; From 5dff649e308d4924053d0144e963cc0865131206 Mon Sep 17 00:00:00 2001 From: castorqin Date: Wed, 18 Oct 2023 16:21:03 +0800 Subject: [PATCH 07/18] [Manager] add default cls manager endpoint,use default cls endpoint to operator cls resource --- .../inlong/manager/service/sink/pulsar/PulsarSinkOperator.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 6a1657b1ec0..6ab38adf5b7 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 @@ -123,7 +123,8 @@ public Map parse2IdParams(StreamSinkEntity streamSink, List Date: Thu, 19 Oct 2023 15:04:36 +0800 Subject: [PATCH 08/18] [Manager] add default cls manager endpoint,use default cls endpoint to operator cls resource --- .../resource/sink/cls/ClsOperator.java | 23 ++++++++----------- 1 file changed, 10 insertions(+), 13 deletions(-) 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 index d0c85a6a53a..d15223b8613 100644 --- 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 @@ -56,7 +56,7 @@ @Service public class ClsOperator { - @Value("${cls.manager.endpoint:cls.internal.tencentcloudapi.com}") + @Value("${cls.manager.endpoint}") private String endpoint; private static final Logger LOG = LoggerFactory.getLogger(ClsOperator.class); private static final String TOPIC_NAME = "topicName"; @@ -64,8 +64,7 @@ public class ClsOperator { private static final long PRECISE_SEARCH = 1L; public String createTopicReturnTopicId(String topicName, String logSetId, String tag, Integer storageDuration, - String secretId, - String secretKey, String region) + String secretId, String secretKey, String region) throws TencentCloudSDKException { ClsClient client = getClsClient(secretId, secretKey, region); CreateTopicRequest req = getCreateTopicRequest(tag, logSetId, topicName, storageDuration); @@ -76,8 +75,8 @@ public String createTopicReturnTopicId(String topicName, String logSetId, String return resp.getTopicId(); } - public void updateTopicTag(String topicId, String tag, String secretId, - String secretKey, String region) throws TencentCloudSDKException { + public void updateTopicTag(String topicId, String tag, String secretId, String secretKey, String region) + throws TencentCloudSDKException { ClsClient client = getClsClient(secretId, secretKey, region); ModifyTopicRequest modifyTopicRequest = new ModifyTopicRequest(); modifyTopicRequest.setTags(convertTags(tag.split(InlongConstants.CENTER_LINE))); @@ -89,8 +88,8 @@ public void updateTopicTag(String topicId, String tag, String secretId, /** * Create topic index by tokenizer */ - public void createTopicIndex(String tokenizer, String topicId, String secretId, String secretKey, - String region) throws BusinessException { + public void createTopicIndex(String tokenizer, String topicId, String secretId, String secretKey, String region) + throws BusinessException { LOG.debug("create topic index start for topicId = {}, tokenizer = {}", topicId, tokenizer); if (StringUtils.isBlank(tokenizer)) { @@ -121,8 +120,8 @@ public void createTopicIndex(String tokenizer, String topicId, String secretId, /** * Describe cls topicId by topic name */ - public String describeTopicIDByTopicName(String topicName, String logSetId, String secretId, - String secretKey, String region) { + public String describeTopicIDByTopicName(String topicName, String logSetId, String secretId, String secretKey, + String region) { ClsClient clsClient = getClsClient(secretId, secretKey, region); Filter[] filters = getDescribeFilters(topicName, logSetId); DescribeTopicsRequest req = new DescribeTopicsRequest(); @@ -158,8 +157,7 @@ public Filter[] getDescribeFilters(String topicName, String logSetId) { /** * Get cls topic index full text */ - public FullTextInfo getTopicIndexFullText(String secretId, String secretKey, String region, - String topicId) { + public FullTextInfo getTopicIndexFullText(String secretId, String secretKey, String region, String topicId) { ClsClient clsClient = getClsClient(secretId, secretKey, region); DescribeIndexRequest req = new DescribeIndexRequest(); @@ -174,8 +172,7 @@ public FullTextInfo getTopicIndexFullText(String secretId, String secretKey, Str } } - public void updateTopicIndex(String tokenizer, String topicId, - String secretId, String secretKey, String region) { + public void updateTopicIndex(String tokenizer, String topicId, String secretId, String secretKey, String region) { ClsClient clsClient = getClsClient(secretId, secretKey, region); RuleInfo ruleInfo = new RuleInfo(); FullTextInfo fullTextInfo = new FullTextInfo(); From 899096ade69095a9f823d7322e46f5f5b14cc6eb Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 10:20:07 +0800 Subject: [PATCH 09/18] [Manager] update code style --- .../manager/service/core/impl/SortClusterServiceImpl.java | 8 ++++---- .../manager/service/sink/pulsar/PulsarSinkOperator.java | 7 +++++-- 2 files changed, 9 insertions(+), 6 deletions(-) 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 3c8dcba0e0f..aa899794371 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 @@ -37,6 +37,7 @@ import com.google.gson.Gson; import org.apache.commons.codec.digest.DigestUtils; +import org.apache.commons.lang3.ObjectUtils; import org.apache.commons.lang3.StringUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -279,7 +280,7 @@ private List> parseIdParams(List streams, StreamSinkOperator operator = sinkOperatorFactory.getInstance(streamSink.getSinkType()); List fields = fieldMap.get(streamSink.getInlongGroupId()); Map params = operator.parse2IdParams(streamSink, fields, dataNodeInfo); - params = setFiledOffset(streamSink, params); + setFiledOffset(streamSink, params); return params; } catch (Exception e) { LOGGER.error("fail to parse id params of groupId={}, streamId={} name={}, type={}}", @@ -292,16 +293,15 @@ private List> parseIdParams(List streams, .collect(Collectors.toList()); } - private Map setFiledOffset(StreamSinkEntity streamSink, Map params) { + private void setFiledOffset(StreamSinkEntity streamSink, Map params) { SortSourceStreamInfo sortSourceStreamInfo = allStreams.get(streamSink.getInlongGroupId()) .get(streamSink.getInlongStreamId()); InlongStreamExtParam inlongStreamExtParam = JsonUtils.parseObject( sortSourceStreamInfo.getExtParams(), InlongStreamExtParam.class); - if (inlongStreamExtParam != null && !inlongStreamExtParam.getUseExtendedFields()) { + if (ObjectUtils.anyNotNull(inlongStreamExtParam) && !inlongStreamExtParam.getUseExtendedFields()) { params.put(FILED_OFFSET, String.valueOf(0)); } - return params; } private Map parseSinkParams(DataNodeInfo nodeInfo) { 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 6ab38adf5b7..dcad49ae908 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 @@ -46,6 +46,8 @@ import java.util.List; import java.util.Map; +import static org.apache.inlong.manager.common.consts.InlongConstants.PULSAR_TOPIC_FORMAT; + /** * Pulsar sink operator */ @@ -123,8 +125,9 @@ public Map parse2IdParams(StreamSinkEntity streamSink, List Date: Mon, 23 Oct 2023 10:39:13 +0800 Subject: [PATCH 10/18] [Manager] Add sink id to parse sink field --- .../main/resources/mappers/StreamSinkFieldEntityMapper.xml | 3 ++- .../inlong/manager/pojo/sort/standalone/SortFieldInfo.java | 1 + .../manager/service/core/impl/SortClusterServiceImpl.java | 7 ++++--- 3 files changed, 7 insertions(+), 4 deletions(-) diff --git a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml index bc4a2bc127b..2fae94247fe 100644 --- a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml +++ b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkFieldEntityMapper.xml @@ -123,8 +123,9 @@ select inlong_group_id, inlong_stream_id, + sink_id, field_name - from inlong_stream_field + from stream_sink_field where is_deleted = 0 order by id asc diff --git a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortFieldInfo.java b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortFieldInfo.java index b4fcbd9fd94..fc8ad991569 100644 --- a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortFieldInfo.java +++ b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortFieldInfo.java @@ -24,5 +24,6 @@ public class SortFieldInfo { private String inlongGroupId; private String inlongStreamId; + private Integer sinkId; private String fieldName; } 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 aa899794371..170b2952493 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 @@ -81,7 +81,8 @@ public class SortClusterServiceImpl implements SortClusterService { private static final String KEY_GROUP_ID = "inlongGroupId"; private static final String KEY_STREAM_ID = "inlongStreamId"; private static final String FILED_OFFSET = "fieldOffset"; - private Map> fieldMap; + // key: sink id, value: fileNames + private Map> fieldMap; // key : sort cluster name, value : md5 private Map sortClusterMd5Map = new ConcurrentHashMap<>(); @@ -176,7 +177,7 @@ private void reloadAllClusterConfig() { List fieldInfos = sortConfigLoader.loadAllFields(); fieldMap = new HashMap<>(); fieldInfos.forEach(info -> { - List fields = fieldMap.computeIfAbsent(info.getInlongGroupId(), k -> new ArrayList<>()); + List fields = fieldMap.computeIfAbsent(info.getSinkId(), k -> new ArrayList<>()); fields.add(info.getFieldName()); }); @@ -278,7 +279,7 @@ private List> parseIdParams(List streams, .map(streamSink -> { try { StreamSinkOperator operator = sinkOperatorFactory.getInstance(streamSink.getSinkType()); - List fields = fieldMap.get(streamSink.getInlongGroupId()); + List fields = fieldMap.get(streamSink.getId()); Map params = operator.parse2IdParams(streamSink, fields, dataNodeInfo); setFiledOffset(streamSink, params); return params; From 0b3e72d06bed4c34c280b9b75e17d1bab6af2807 Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 11:15:51 +0800 Subject: [PATCH 11/18] [Manager] Add sink id to parse sink field --- .../manager-web/src/main/resources/application-dev.properties | 3 +++ 1 file changed, 3 insertions(+) diff --git a/inlong-manager/manager-web/src/main/resources/application-dev.properties b/inlong-manager/manager-web/src/main/resources/application-dev.properties index 376d9b9f73e..25eee3ef5dc 100644 --- a/inlong-manager/manager-web/src/main/resources/application-dev.properties +++ b/inlong-manager/manager-web/src/main/resources/application-dev.properties @@ -109,3 +109,6 @@ group.deleted.batchSize=100 group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 + +# when operator cls resource,need config cls.manager.endpoint +cls.manager.endpoint=cls.tencentcloudapi.com From a2d1f0ab21f9c44c4630bc20e35caab87f6af9d8 Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 11:34:43 +0800 Subject: [PATCH 12/18] [Manager] Add sink id to parse sink field --- .../manager-web/src/main/resources/application-dev.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/inlong-manager/manager-web/src/main/resources/application-dev.properties b/inlong-manager/manager-web/src/main/resources/application-dev.properties index 25eee3ef5dc..44a7c018e95 100644 --- a/inlong-manager/manager-web/src/main/resources/application-dev.properties +++ b/inlong-manager/manager-web/src/main/resources/application-dev.properties @@ -110,5 +110,5 @@ group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 -# when operator cls resource,need config cls.manager.endpoint +# tencent cloud log service endpoint,The Operator cls resource by it cls.manager.endpoint=cls.tencentcloudapi.com From 11583ebdb7de721781349d68b1209d803c42990d Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 11:41:07 +0800 Subject: [PATCH 13/18] [Manager] Add sink id to parse sink field --- .../manager-web/src/main/resources/application-dev.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/inlong-manager/manager-web/src/main/resources/application-dev.properties b/inlong-manager/manager-web/src/main/resources/application-dev.properties index 44a7c018e95..729bbc2f0e7 100644 --- a/inlong-manager/manager-web/src/main/resources/application-dev.properties +++ b/inlong-manager/manager-web/src/main/resources/application-dev.properties @@ -110,5 +110,5 @@ group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 -# tencent cloud log service endpoint,The Operator cls resource by it +# tencent cloud log service endpoint, The Operator cls resource by it cls.manager.endpoint=cls.tencentcloudapi.com From 0888c867610da60b7567e5529f5aa6ae61915620 Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 11:45:03 +0800 Subject: [PATCH 14/18] [Manager] Add sink id to parse sink field --- .../manager-web/src/main/resources/application-dev.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/inlong-manager/manager-web/src/main/resources/application-dev.properties b/inlong-manager/manager-web/src/main/resources/application-dev.properties index 729bbc2f0e7..27cf7c7d311 100644 --- a/inlong-manager/manager-web/src/main/resources/application-dev.properties +++ b/inlong-manager/manager-web/src/main/resources/application-dev.properties @@ -111,4 +111,4 @@ group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 # tencent cloud log service endpoint, The Operator cls resource by it -cls.manager.endpoint=cls.tencentcloudapi.com +cls.manager.endpoint=xxx From 381280c790c027a97fe8c9433f6976f50600d438 Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 14:34:09 +0800 Subject: [PATCH 15/18] [Manager] Add sink id to parse sink field --- .../manager-web/src/main/resources/application-test.properties | 3 +++ 1 file changed, 3 insertions(+) diff --git a/inlong-manager/manager-web/src/main/resources/application-test.properties b/inlong-manager/manager-web/src/main/resources/application-test.properties index 376d9b9f73e..27cf7c7d311 100644 --- a/inlong-manager/manager-web/src/main/resources/application-test.properties +++ b/inlong-manager/manager-web/src/main/resources/application-test.properties @@ -109,3 +109,6 @@ group.deleted.batchSize=100 group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 + +# tencent cloud log service endpoint, The Operator cls resource by it +cls.manager.endpoint=xxx From 558f8aab0bda8340a83f049a26ce0f596b9cd430 Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 17:10:08 +0800 Subject: [PATCH 16/18] [Manager] Add sink id to parse sink field --- .../manager-web/src/main/resources/application-prod.properties | 3 +++ .../manager-web/src/main/resources/application.properties | 3 +++ 2 files changed, 6 insertions(+) diff --git a/inlong-manager/manager-web/src/main/resources/application-prod.properties b/inlong-manager/manager-web/src/main/resources/application-prod.properties index c47fc92334b..dcf92021155 100644 --- a/inlong-manager/manager-web/src/main/resources/application-prod.properties +++ b/inlong-manager/manager-web/src/main/resources/application-prod.properties @@ -108,3 +108,6 @@ group.deleted.batchSize=100 group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 + +# tencent cloud log service endpoint, The Operator cls resource by it +cls.manager.endpoint=xxx diff --git a/inlong-manager/manager-web/src/main/resources/application.properties b/inlong-manager/manager-web/src/main/resources/application.properties index c434645b2fc..e1c6fc76d34 100644 --- a/inlong-manager/manager-web/src/main/resources/application.properties +++ b/inlong-manager/manager-web/src/main/resources/application.properties @@ -63,3 +63,6 @@ openapi.auth.enabled=false # Audit view by role, see audit id definitions: https://inlong.apache.org/docs/modules/audit/overview#audit-id audit.admin.ids=3,4,5,6 audit.user.ids=3,4,5,6 + +# tencent cloud log service endpoint, The Operator cls resource by it +cls.manager.endpoint=xxx From df20788e5f0cea03d02de72146a379de38107631 Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 18:24:56 +0800 Subject: [PATCH 17/18] [Manager] Add sink id to parse sink field --- .../manager-web/src/main/resources/application-dev.properties | 2 +- .../manager-web/src/main/resources/application-prod.properties | 2 +- .../manager-web/src/main/resources/application-test.properties | 2 +- .../manager-web/src/main/resources/application.properties | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) diff --git a/inlong-manager/manager-web/src/main/resources/application-dev.properties b/inlong-manager/manager-web/src/main/resources/application-dev.properties index 27cf7c7d311..0bcebdc06e3 100644 --- a/inlong-manager/manager-web/src/main/resources/application-dev.properties +++ b/inlong-manager/manager-web/src/main/resources/application-dev.properties @@ -111,4 +111,4 @@ group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 # tencent cloud log service endpoint, The Operator cls resource by it -cls.manager.endpoint=xxx +cls.manager.endpoint=127.0.0.1 diff --git a/inlong-manager/manager-web/src/main/resources/application-prod.properties b/inlong-manager/manager-web/src/main/resources/application-prod.properties index dcf92021155..a9c55b39b39 100644 --- a/inlong-manager/manager-web/src/main/resources/application-prod.properties +++ b/inlong-manager/manager-web/src/main/resources/application-prod.properties @@ -110,4 +110,4 @@ group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 # tencent cloud log service endpoint, The Operator cls resource by it -cls.manager.endpoint=xxx +cls.manager.endpoint=127.0.0.1 diff --git a/inlong-manager/manager-web/src/main/resources/application-test.properties b/inlong-manager/manager-web/src/main/resources/application-test.properties index 27cf7c7d311..0bcebdc06e3 100644 --- a/inlong-manager/manager-web/src/main/resources/application-test.properties +++ b/inlong-manager/manager-web/src/main/resources/application-test.properties @@ -111,4 +111,4 @@ group.deleted.enabled=false metrics.audit.proxy.hosts=127.0.0.1:10081 # tencent cloud log service endpoint, The Operator cls resource by it -cls.manager.endpoint=xxx +cls.manager.endpoint=127.0.0.1 diff --git a/inlong-manager/manager-web/src/main/resources/application.properties b/inlong-manager/manager-web/src/main/resources/application.properties index e1c6fc76d34..6b56dfb3d99 100644 --- a/inlong-manager/manager-web/src/main/resources/application.properties +++ b/inlong-manager/manager-web/src/main/resources/application.properties @@ -65,4 +65,4 @@ audit.admin.ids=3,4,5,6 audit.user.ids=3,4,5,6 # tencent cloud log service endpoint, The Operator cls resource by it -cls.manager.endpoint=xxx +cls.manager.endpoint=127.0.0.1 From 9ae6c91652db391baf95d0277a34b055d1b25b0c Mon Sep 17 00:00:00 2001 From: castorqin Date: Mon, 23 Oct 2023 19:33:15 +0800 Subject: [PATCH 18/18] [Manager] Add sink id to parse sink field --- .../src/main/resources/application-unit-test.properties | 3 +++ 1 file changed, 3 insertions(+) diff --git a/inlong-manager/manager-test/src/main/resources/application-unit-test.properties b/inlong-manager/manager-test/src/main/resources/application-unit-test.properties index 862eba47064..aa9155e7353 100644 --- a/inlong-manager/manager-test/src/main/resources/application-unit-test.properties +++ b/inlong-manager/manager-test/src/main/resources/application-unit-test.properties @@ -67,3 +67,6 @@ common.http-client.validateAfterInactivity=5000 common.http-client.connectionTimeout=3000 common.http-client.readTimeout=10000 common.http-client.connectionRequestTimeout=3000 + +# tencent cloud log service endpoint, The Operator cls resource by it +cls.manager.endpoint=127.0.0.1