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 e673c2565 [INLONG-6104][Manager] Support getting backup info in
getSortSource (#6275)
e673c2565 is described below
commit e673c2565bf24e2550444942bab5e165bea60d81
Author: vernedeng <[email protected]>
AuthorDate: Sat Oct 29 19:51:13 2022 +0800
[INLONG-6104][Manager] Support getting backup info in getSortSource (#6275)
---
.../dao/mapper/InlongClusterEntityMapper.java | 6 +-
.../dao/mapper/InlongGroupEntityMapper.java | 6 +-
.../dao/mapper/InlongGroupExtEntityMapper.java | 10 +-
.../dao/mapper/InlongStreamEntityMapper.java | 8 +-
.../dao/mapper/InlongStreamExtEntityMapper.java | 6 +
.../manager/dao/mapper/StreamSinkEntityMapper.java | 10 +-
.../mappers/InlongClusterEntityMapper.xml | 3 +-
.../resources/mappers/InlongGroupEntityMapper.xml | 2 +-
.../mappers/InlongGroupExtEntityMapper.xml | 7 +
.../resources/mappers/InlongStreamEntityMapper.xml | 8 +
.../mappers/InlongStreamExtEntityMapper.xml | 20 +-
.../resources/mappers/StreamSinkEntityMapper.xml | 8 +-
.../sort/standalone/SortSourceClusterInfo.java | 5 +-
.../pojo/sort/standalone/SortSourceGroupInfo.java | 6 +-
.../pojo/sort/standalone/SortSourceStreamInfo.java | 34 +--
...reamInfo.java => SortSourceStreamSinkInfo.java} | 8 +-
.../manager/service/core/SortConfigLoader.java | 70 ++++++
.../service/core/impl/SortConfigLoaderImpl.java | 109 ++++++++
.../service/core/impl/SortSourceServiceImpl.java | 280 ++++++++++-----------
.../service/group/InlongGroupOperator4Pulsar.java | 17 +-
.../service/group/InlongGroupServiceImpl.java | 2 +-
.../manager/service/sort/SortServiceImplTest.java | 113 ++++++---
22 files changed, 487 insertions(+), 251 deletions(-)
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
index db7ae4a70..1f993527e 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongClusterEntityMapper.java
@@ -17,7 +17,10 @@
package org.apache.inlong.manager.dao.mapper;
+import org.apache.ibatis.annotations.Options;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
import org.apache.inlong.manager.dao.entity.InlongClusterEntity;
import org.apache.inlong.manager.pojo.cluster.ClusterPageRequest;
import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
@@ -49,7 +52,8 @@ public interface InlongClusterEntityMapper {
*
* @return All cluster info.
*/
- List<SortSourceClusterInfo> selectAllClusters();
+ @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize =
Integer.MIN_VALUE)
+ Cursor<SortSourceClusterInfo> selectAllClusters();
List<InlongClusterEntity> selectByClusterTag(String clusterTag);
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
index b9ba70058..7dc92543d 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupEntityMapper.java
@@ -17,7 +17,10 @@
package org.apache.inlong.manager.dao.mapper;
+import org.apache.ibatis.annotations.Options;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
import org.apache.inlong.manager.dao.entity.InlongGroupEntity;
import org.apache.inlong.manager.pojo.group.InlongGroupBriefInfo;
import org.apache.inlong.manager.pojo.group.InlongGroupPageRequest;
@@ -52,7 +55,8 @@ public interface InlongGroupEntityMapper {
*
* @return All inlong group info.
*/
- List<SortSourceGroupInfo> selectAllGroups();
+ @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize =
Integer.MIN_VALUE)
+ Cursor<SortSourceGroupInfo> selectAllGroups();
/**
* Select all groups which are logical deleted before the specified last
modify time
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
index 4e7c1d4ab..38b694f1c 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongGroupExtEntityMapper.java
@@ -17,7 +17,10 @@
package org.apache.inlong.manager.dao.mapper;
+import org.apache.ibatis.annotations.Options;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
import org.springframework.stereotype.Repository;
@@ -38,7 +41,12 @@ public interface InlongGroupExtEntityMapper {
int updateByPrimaryKey(InlongGroupExtEntity record);
- InlongGroupExtEntity selectByUniqueKey(@Param("groupId") String groupId,
@Param("keyName") String keyName);
+ InlongGroupExtEntity selectByUniqueKey(
+ @Param("inlongGroupId") String inlongGroupId,
+ @Param("keyName") String keyName);
+
+ @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize =
Integer.MIN_VALUE)
+ Cursor<InlongGroupExtEntity> selectByKeyName(@Param("keyName") String
keyName);
/**
* Insert data in batches
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
index 66f101a27..b676dcfeb 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamEntityMapper.java
@@ -17,8 +17,12 @@
package org.apache.inlong.manager.dao.mapper;
+import org.apache.ibatis.annotations.Options;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
import org.apache.inlong.manager.dao.entity.InlongStreamEntity;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamBriefInfo;
import org.apache.inlong.manager.pojo.stream.InlongStreamPageRequest;
import org.springframework.stereotype.Repository;
@@ -50,6 +54,9 @@ public interface InlongStreamEntityMapper {
int selectCountByGroupId(@Param("groupId") String groupId);
+ @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize =
Integer.MIN_VALUE)
+ Cursor<SortSourceStreamInfo> selectAllStreams();
+
int updateByPrimaryKey(InlongStreamEntity record);
int updateByIdentifierSelective(InlongStreamEntity streamEntity);
@@ -70,5 +77,4 @@ public interface InlongStreamEntityMapper {
* @return rows deleted
*/
int deleteByInlongGroupIds(@Param("groupIdList") List<String> groupIdList);
-
}
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
index 2c7316bf6..b9457210b 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/InlongStreamExtEntityMapper.java
@@ -17,7 +17,10 @@
package org.apache.inlong.manager.dao.mapper;
+import org.apache.ibatis.annotations.Options;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
import org.springframework.stereotype.Repository;
@@ -47,6 +50,9 @@ public interface InlongStreamExtEntityMapper {
InlongStreamExtEntity selectByKey(@Param("groupId") String groupId,
@Param("streamId") String streamId,
@Param("keyName") String keyName);
+ @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize =
Integer.MIN_VALUE)
+ Cursor<InlongStreamExtEntity> selectByKeyName(@Param("keyName") String
keyName);
+
int updateByPrimaryKey(InlongStreamExtEntity record);
/**
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
index 7142976dc..20fbd9e55 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/mapper/StreamSinkEntityMapper.java
@@ -17,13 +17,16 @@
package org.apache.inlong.manager.dao.mapper;
+import org.apache.ibatis.annotations.Options;
import org.apache.ibatis.annotations.Param;
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.ibatis.mapping.ResultSetType;
import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
import org.apache.inlong.manager.pojo.sink.SinkBriefInfo;
import org.apache.inlong.manager.pojo.sink.SinkInfo;
import org.apache.inlong.manager.pojo.sink.SinkPageRequest;
import org.apache.inlong.manager.pojo.sort.standalone.SortIdInfo;
-import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
import org.apache.inlong.manager.pojo.sort.standalone.SortTaskInfo;
import org.springframework.stereotype.Repository;
@@ -117,6 +120,8 @@ public interface StreamSinkEntityMapper {
*/
List<SinkInfo> selectAllConfig(@Param("groupId") String groupId,
@Param("idList") List<String> streamIdList);
+ List<StreamSinkEntity> selectAllStreamSinks();
+
/**
* Select all tasks for sort-standalone
*
@@ -136,7 +141,8 @@ public interface StreamSinkEntityMapper {
*
* @return All stream info
*/
- List<SortSourceStreamInfo> selectAllStreams();
+ @Options(resultSetType = ResultSetType.FORWARD_ONLY, fetchSize =
Integer.MIN_VALUE)
+ Cursor<SortSourceStreamSinkInfo> selectAllStreams();
int updateByIdSelective(StreamSinkEntity record);
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
index d9d7fef62..931ca5f83 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongClusterEntityMapper.xml
@@ -164,9 +164,10 @@
</where>
order by modify_time desc
</select>
- <select id="selectAllClusters"
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo">
+ <select id="selectAllClusters"
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo"
>
select name,
type,
+ url,
cluster_tags,
ext_tag,
ext_params
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
index e3b832339..47a97759f 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupEntityMapper.xml
@@ -202,7 +202,7 @@
<select id="selectAllGroups"
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo">
select inlong_group_id as groupId,
inlong_cluster_tag as clusterTag,
- mq_resource as topic,
+ mq_resource as mqResource,
ext_params as extParams,
mq_type as mqType
from inlong_group
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
index dba70b312..9c16c25d7 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongGroupExtEntityMapper.xml
@@ -53,6 +53,13 @@
and key_name = #{keyName, jdbcType=VARCHAR}
and is_deleted = 0
</select>
+ <select id="selectByKeyName"
resultType="org.apache.inlong.manager.dao.entity.InlongGroupExtEntity">
+ select
+ <include refid="Base_Column_List"/>
+ from inlong_group_ext
+ where key_name = #{keyName, jdbcType=VARCHAR}
+ and is_deleted = 0
+ </select>
<delete id="deleteByPrimaryKey" parameterType="java.lang.Integer">
delete
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
index 136c33f12..231769ccc 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamEntityMapper.xml
@@ -285,6 +285,14 @@
where inlong_group_id = #{groupId, jdbcType=VARCHAR}
and is_deleted = 0
</select>
+ <select id="selectAllStreams"
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo">
+ select
+ inlong_group_id,
+ inlong_stream_id,
+ mq_resource
+ from inlong_stream
+ where is_deleted = 0
+ </select>
<select id="selectCountByGroupId" resultType="java.lang.Integer">
select count(1)
from inlong_stream
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
index 25e4e1a7a..03c785f58 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/InlongStreamExtEntityMapper.xml
@@ -35,30 +35,29 @@
<insert id="insert"
parameterType="org.apache.inlong.manager.dao.entity.InlongStreamExtEntity">
insert into inlong_stream_ext (id, inlong_group_id, inlong_stream_id,
- key_name, is_deleted, modify_time,
+ key_name, modify_time,
key_value)
values (#{id,jdbcType=INTEGER}, #{inlongGroupId,jdbcType=VARCHAR},
#{inlongStreamId,jdbcType=VARCHAR},
- #{keyName,jdbcType=VARCHAR}, #{isDeleted,jdbcType=INTEGER},
#{modifyTime,jdbcType=TIMESTAMP},
+ #{keyName,jdbcType=VARCHAR}, #{modifyTime,jdbcType=TIMESTAMP},
#{keyValue,jdbcType=LONGVARCHAR})
</insert>
<!-- Bulk insert-->
<insert id="insertAll" parameterType="java.util.List">
insert into inlong_stream_ext
- (id, inlong_group_id, inlong_stream_id, key_name, key_value,
is_deleted)
+ (id, inlong_group_id, inlong_stream_id, key_name, key_value)
values
<foreach collection="extList" separator="," index="index" item="item">
- (#{item.id}, #{item.inlongGroupId}, #{item.inlongStreamId},
#{item.keyName}, #{item.keyValue},
- #{item.isDeleted})
+ (#{item.id}, #{item.inlongGroupId}, #{item.inlongStreamId},
#{item.keyName}, #{item.keyValue})
</foreach>
</insert>
<!-- Bulk insert, update if exists-->
<insert id="insertOnDuplicateKeyUpdate" parameterType="java.util.List">
insert into inlong_stream_ext
- (id, inlong_group_id, inlong_stream_id, key_name, key_value,
is_deleted)
+ (id, inlong_group_id, inlong_stream_id, key_name, key_value)
values
<foreach collection="extList" separator="," index="index" item="item">
(#{item.id}, #{item.inlongGroupId}, #{item.inlongStreamId},
- #{item.keyName}, #{item.keyValue}, #{item.isDeleted})
+ #{item.keyName}, #{item.keyValue})
</foreach>
ON DUPLICATE KEY UPDATE
inlong_group_id = VALUES(inlong_group_id),
@@ -87,6 +86,13 @@
and key_name = #{keyName, jdbcType=VARCHAR}
and is_deleted = 0
</select>
+ <select id="selectByKeyName"
resultType="org.apache.inlong.manager.dao.entity.InlongStreamExtEntity">
+ select
+ <include refid="Base_Column_List"/>
+ from inlong_stream_ext
+ where key_name = #{keyName, jdbcType=VARCHAR}
+ and is_deleted = 0
+ </select>
<update id="updateByPrimaryKey"
parameterType="org.apache.inlong.manager.dao.entity.InlongStreamExtEntity">
update inlong_stream_ext
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
index c4af01696..aac142d59 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/StreamSinkEntityMapper.xml
@@ -350,6 +350,12 @@
</if>
</where>
</select>
+ <select id="selectAllStreamSinks"
resultType="org.apache.inlong.manager.dao.entity.StreamSinkEntity">
+ select
+ <include refid="Base_Column_List"/>
+ from stream_sink
+ where is_deleted = 0
+ </select>
<select id="selectAllTasks"
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortTaskInfo">
select inlong_cluster_name as sortClusterName,
sort_task_name,
@@ -368,7 +374,7 @@
from stream_sink
where is_deleted = 0
</select>
- <select id="selectAllStreams"
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo">
+ <select id="selectAllStreams"
resultType="org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo">
select inlong_cluster_name as sortClusterName,
sort_task_name,
inlong_group_id as groupId,
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
index 148320178..933982729 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceClusterInfo.java
@@ -31,6 +31,7 @@ import java.util.concurrent.ConcurrentHashMap;
public class SortSourceClusterInfo {
private static final Logger LOGGER =
LoggerFactory.getLogger(SortSourceClusterInfo.class);
+ private static final Gson GSON = new Gson();
private static final Splitter.MapSplitter MAP_SPLITTER =
Splitter.on("&").trimResults()
.withKeyValueSeparator("=");
private static final String KEY_IS_CONSUMABLE = "consumer";
@@ -38,6 +39,7 @@ public class SortSourceClusterInfo {
private static final long serialVersionUID = 1L;
String name;
String type;
+ String url;
String clusterTags;
String extTag;
String extParams;
@@ -47,8 +49,7 @@ public class SortSourceClusterInfo {
public Map<String, String> getExtParamsMap() {
if (extParamsMap.isEmpty() && extParams != null) {
try {
- Gson gson = new Gson();
- extParamsMap = gson.fromJson(extParams, Map.class);
+ extParamsMap = GSON.fromJson(extParams, Map.class);
} catch (Throwable t) {
LOGGER.error("fail to parse cluster ext params", t);
}
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
index 8d887d609..b7b951652 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceGroupInfo.java
@@ -30,13 +30,14 @@ import java.util.concurrent.ConcurrentHashMap;
public class SortSourceGroupInfo {
private static final Logger LOGGER =
LoggerFactory.getLogger(SortSourceGroupInfo.class);
+ private static final Gson GSON = new Gson();
private static final String KEY_BACKUP_CLUSTER_TAG = "backup_cluster_tag";
private static final String KEY_BACKUP_TOPIC = "backup_topic";
private static final long serialVersionUID = 1L;
String groupId;
String clusterTag;
- String topic;
+ String mqResource;
String extParams;
String mqType;
Map<String, String> extParamsMap = new ConcurrentHashMap<>();
@@ -44,8 +45,7 @@ public class SortSourceGroupInfo {
public Map<String, String> getExtParamsMap() {
if (extParamsMap.isEmpty() && StringUtils.isNotBlank(extParams)) {
try {
- Gson gson = new Gson();
- extParamsMap = gson.fromJson(extParams, Map.class);
+ extParamsMap = GSON.fromJson(extParams, Map.class);
} catch (Throwable t) {
LOGGER.error("fail to parse group ext params", t);
}
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
index 6462e3e77..bd89a573c 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
@@ -17,39 +17,11 @@
package org.apache.inlong.manager.pojo.sort.standalone;
-import com.google.gson.Gson;
import lombok.Data;
-import org.apache.commons.lang3.StringUtils;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
@Data
public class SortSourceStreamInfo {
-
- private static final long serialVersionUID = 1L;
- private static final Logger LOGGER =
LoggerFactory.getLogger(SortSourceStreamInfo.class);
- String sortClusterName;
- String sortTaskName;
- String groupId;
- String extParams;
- Map<String, String> extParamsMap;
-
- public Map<String, String> getExtParamsMap() {
- if (extParamsMap != null) {
- return extParamsMap;
- }
- if (StringUtils.isNotBlank(extParams)) {
- try {
- Gson gson = new Gson();
- extParamsMap = gson.fromJson(extParams, Map.class);
- } catch (Throwable t) {
- LOGGER.error("fail to parse source stream ext params", t);
- extParamsMap = new ConcurrentHashMap<>();
- }
- }
- return extParamsMap;
- }
+ private String inlongGroupId;
+ private String inlongStreamId;
+ private String mqResource;
}
diff --git
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamSinkInfo.java
similarity index 90%
copy from
inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
copy to
inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamSinkInfo.java
index 6462e3e77..4aa2e1168 100644
---
a/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamInfo.java
+++
b/inlong-manager/manager-pojo/src/main/java/org/apache/inlong/manager/pojo/sort/standalone/SortSourceStreamSinkInfo.java
@@ -27,10 +27,11 @@ import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Data
-public class SortSourceStreamInfo {
+public class SortSourceStreamSinkInfo {
private static final long serialVersionUID = 1L;
- private static final Logger LOGGER =
LoggerFactory.getLogger(SortSourceStreamInfo.class);
+ private static final Logger LOGGER =
LoggerFactory.getLogger(SortSourceStreamSinkInfo.class);
+ private static final Gson GSON = new Gson();
String sortClusterName;
String sortTaskName;
String groupId;
@@ -43,8 +44,7 @@ public class SortSourceStreamInfo {
}
if (StringUtils.isNotBlank(extParams)) {
try {
- Gson gson = new Gson();
- extParamsMap = gson.fromJson(extParams, Map.class);
+ extParamsMap = GSON.fromJson(extParams, Map.class);
} catch (Throwable t) {
LOGGER.error("fail to parse source stream ext params", t);
extParamsMap = new ConcurrentHashMap<>();
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/SortConfigLoader.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/SortConfigLoader.java
new file mode 100644
index 000000000..7bfa2a8ec
--- /dev/null
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/SortConfigLoader.java
@@ -0,0 +1,70 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.manager.service.core;
+
+import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
+import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
+
+import java.util.List;
+
+/**
+ * Loader for sort service to load configs thought Cursor
+ */
+public interface SortConfigLoader {
+ /**
+ * Load all clusters by cursor
+ * @return List of clusters, including MQ cluster and DataNode cluster.
+ */
+ List<SortSourceClusterInfo> loadAllClusters();
+
+ /**
+ * Load stream sinks by cursor
+ * @return List of Stream sinks.
+ */
+ List<SortSourceStreamSinkInfo> loadAllStreamSinks();
+
+ /**
+ * Load groups by cursor
+ * @return List of group info
+ */
+ List<SortSourceGroupInfo> loadAllGroup();
+
+ /**
+ * Load group backup info by cursor
+ * @param keyName Key name
+ * @return List of group backup info
+ */
+ List<InlongGroupExtEntity> loadGroupBackupInfo(String keyName);
+
+ /**
+ * Load stream backup info by cursor
+ * @param keyName Key name
+ * @return List of stream backup info
+ */
+ List<InlongStreamExtEntity> loadStreamBackupInfo(String keyName);
+
+ /**
+ * Load all inlong stream info by cursor
+ * @return List of stream info
+ */
+ List<SortSourceStreamInfo> loadAllStreams();
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortConfigLoaderImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortConfigLoaderImpl.java
new file mode 100644
index 000000000..f442f9a46
--- /dev/null
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortConfigLoaderImpl.java
@@ -0,0 +1,109 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.manager.service.core.impl;
+
+import org.apache.ibatis.cursor.Cursor;
+import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
+import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
+import org.apache.inlong.manager.dao.mapper.InlongClusterEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongGroupEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongGroupExtEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongStreamEntityMapper;
+import org.apache.inlong.manager.dao.mapper.InlongStreamExtEntityMapper;
+import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
+import org.apache.inlong.manager.service.core.SortConfigLoader;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+import org.springframework.transaction.annotation.Transactional;
+
+import java.util.ArrayList;
+import java.util.List;
+
+@Service
+public class SortConfigLoaderImpl implements SortConfigLoader {
+ @Autowired
+ private InlongClusterEntityMapper clusterEntityMapper;
+ @Autowired
+ private StreamSinkEntityMapper streamSinkEntityMapper;
+ @Autowired
+ private InlongGroupEntityMapper inlongGroupEntityMapper;
+ @Autowired
+ private InlongGroupExtEntityMapper inlongGroupExtEntityMapper;
+ @Autowired
+ private InlongStreamExtEntityMapper inlongStreamExtEntityMapper;
+ @Autowired
+ private InlongStreamEntityMapper inlongStreamEntityMapper;
+
+ @Transactional
+ @Override
+ public List<SortSourceClusterInfo> loadAllClusters() {
+ Cursor<SortSourceClusterInfo> cursor =
clusterEntityMapper.selectAllClusters();
+ List<SortSourceClusterInfo> allClusters = new ArrayList<>();
+ cursor.forEach(allClusters::add);
+ return allClusters;
+ }
+
+ @Transactional
+ @Override
+ public List<SortSourceStreamSinkInfo> loadAllStreamSinks() {
+ Cursor<SortSourceStreamSinkInfo> cursor =
streamSinkEntityMapper.selectAllStreams();
+ List<SortSourceStreamSinkInfo> allStreamSinks = new ArrayList<>();
+ cursor.forEach(allStreamSinks::add);
+ return allStreamSinks;
+ }
+
+ @Transactional
+ @Override
+ public List<SortSourceGroupInfo> loadAllGroup() {
+ Cursor<SortSourceGroupInfo> cursor =
inlongGroupEntityMapper.selectAllGroups();
+ List<SortSourceGroupInfo> allGroups = new ArrayList<>();
+ cursor.forEach(allGroups::add);
+ return allGroups;
+ }
+
+ @Transactional
+ @Override
+ public List<InlongGroupExtEntity> loadGroupBackupInfo(String keyName) {
+ Cursor<InlongGroupExtEntity> cursor =
inlongGroupExtEntityMapper.selectByKeyName(keyName);
+ List<InlongGroupExtEntity> groupBackupInfo = new ArrayList<>();
+ cursor.forEach(groupBackupInfo::add);
+ return groupBackupInfo;
+ }
+
+ @Transactional
+ @Override
+ public List<InlongStreamExtEntity> loadStreamBackupInfo(String keyName) {
+ Cursor<InlongStreamExtEntity> cursor =
inlongStreamExtEntityMapper.selectByKeyName(keyName);
+ List<InlongStreamExtEntity> streamBackupInfo = new ArrayList<>();
+ cursor.forEach(streamBackupInfo::add);
+ return streamBackupInfo;
+ }
+
+ @Transactional
+ @Override
+ public List<SortSourceStreamInfo> loadAllStreams() {
+ Cursor<SortSourceStreamInfo> cursor =
inlongStreamEntityMapper.selectAllStreams();
+ List<SortSourceStreamInfo> allStreams = new ArrayList<>();
+ cursor.forEach(allStreams::add);
+ return allStreams;
+ }
+}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
index 611e8d02e..cadf713c5 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/SortSourceServiceImpl.java
@@ -20,17 +20,21 @@ package org.apache.inlong.manager.service.core.impl;
import com.google.gson.Gson;
import org.apache.commons.codec.digest.DigestUtils;
import org.apache.commons.lang3.StringUtils;
+import org.apache.inlong.common.constant.ClusterSwitch;
import org.apache.inlong.common.pojo.sdk.CacheZone;
import org.apache.inlong.common.pojo.sdk.CacheZoneConfig;
import org.apache.inlong.common.pojo.sdk.SortSourceConfigResponse;
import org.apache.inlong.common.pojo.sdk.Topic;
import org.apache.inlong.manager.common.consts.MQType;
-import org.apache.inlong.manager.dao.mapper.InlongClusterEntityMapper;
-import org.apache.inlong.manager.dao.mapper.InlongGroupEntityMapper;
-import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
+import org.apache.inlong.manager.common.enums.ClusterType;
+import org.apache.inlong.manager.common.exceptions.BusinessException;
+import org.apache.inlong.manager.dao.entity.InlongGroupExtEntity;
+import org.apache.inlong.manager.dao.entity.InlongStreamExtEntity;
import org.apache.inlong.manager.pojo.sort.standalone.SortSourceClusterInfo;
import org.apache.inlong.manager.pojo.sort.standalone.SortSourceGroupInfo;
import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamInfo;
+import org.apache.inlong.manager.pojo.sort.standalone.SortSourceStreamSinkInfo;
+import org.apache.inlong.manager.service.core.SortConfigLoader;
import org.apache.inlong.manager.service.core.SortSourceService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -45,9 +49,7 @@ import java.util.Collection;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
-import java.util.Locale;
import java.util.Map;
-import java.util.Objects;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
@@ -74,10 +76,8 @@ public class SortSourceServiceImpl implements
SortSourceService {
add(MQType.PULSAR);
}
};
- private static final String KEY_SERVICE_URL = "serviceUrl";
private static final String KEY_AUTH = "authentication";
private static final String KEY_TENANT = "tenant";
- private static final String KEY_NAME_SPACE = "namespace";
private static final int RESPONSE_CODE_SUCCESS = 0;
private static final int RESPONSE_CODE_NO_UPDATE = 1;
@@ -93,12 +93,16 @@ public class SortSourceServiceImpl implements
SortSourceService {
*/
private Map<String, Map<String, CacheZoneConfig>> sortSourceConfigMap =
new ConcurrentHashMap<>();
+ private Map<String, List<SortSourceClusterInfo>> mqClusters;
+ private Map<String, SortSourceGroupInfo> groupInfos;
+ private Map<String, SortSourceStreamInfo> allStreams;
+ private Map<String, String> backupClusterTag;
+ private Map<String, String> backupGroupMqResource;
+ private Map<String, String> backupStreamMqResource;
+ private Map<String, Map<String, List<String>>> groupMap;
+
@Autowired
- private InlongClusterEntityMapper clusterEntityMapper;
- @Autowired
- private StreamSinkEntityMapper streamSinkEntityMapper;
- @Autowired
- private InlongGroupEntityMapper inlongGroupEntityMapper;
+ private SortConfigLoader configLoader;
@PostConstruct
public void initialize() {
@@ -115,7 +119,8 @@ public class SortSourceServiceImpl implements
SortSourceService {
public void reload() {
LOGGER.debug("start to reload sort config.");
try {
- reloadAllSourceConfig();
+ reloadAllConfigs();
+ parseAll();
} catch (Throwable t) {
LOGGER.error("fail to reload all source config", t);
}
@@ -177,16 +182,21 @@ public class SortSourceServiceImpl implements
SortSourceService {
}
- private void reloadAllSourceConfig() {
+ private void reloadAllConfigs() {
- // get all streams.
- List<SortSourceStreamInfo> allStreamInfos =
streamSinkEntityMapper.selectAllStreams().stream()
- .filter(dto -> dto.getSortClusterName() != null &&
dto.getSortTaskName() != null)
- .collect(Collectors.toList());
+ // reload mq cluster and sort cluster
+ List<SortSourceClusterInfo> allClusters =
configLoader.loadAllClusters();
- // convert to Map<clusterName, Map<taskName, List<groupId>>> format.
- Map<String, Map<String, List<String>>> groupMap = new
ConcurrentHashMap<>();
- allStreamInfos.forEach(stream -> {
+ // group mq clusters by cluster tag
+ mqClusters = allClusters.stream()
+ .filter(cluster ->
SUPPORTED_MQ_TYPE.contains(cluster.getType()))
+ .filter(SortSourceClusterInfo::isConsumable)
+
.collect(Collectors.groupingBy(SortSourceClusterInfo::getClusterTags));
+
+ // reload all stream sinks, to Map<clusterName, Map<taskName,
List<groupId>>> format
+ List<SortSourceStreamSinkInfo> allStreamSinks =
configLoader.loadAllStreamSinks();
+ groupMap = new HashMap<>();
+ allStreamSinks.forEach(stream -> {
Map<String, List<String>> task2groupsMap =
groupMap.computeIfAbsent(stream.getSortClusterName(), k ->
new ConcurrentHashMap<>());
List<String> groupIdList =
@@ -194,109 +204,96 @@ public class SortSourceServiceImpl implements
SortSourceService {
groupIdList.add(stream.getGroupId());
});
- // get all groups. group by group id.
- List<SortSourceGroupInfo> groupInfos =
inlongGroupEntityMapper.selectAllGroups();
- Map<String, SortSourceGroupInfo> allId2GroupInfos = groupInfos.stream()
- .filter(dto -> dto.getGroupId() != null)
- .collect(Collectors.toMap(SortSourceGroupInfo::getGroupId, dto
-> dto, (g1, g2) -> g1));
-
- // get all clusters. filter by type and check if consumable, then
group by cluster tag.
- List<SortSourceClusterInfo> clusterInfos =
clusterEntityMapper.selectAllClusters();
- Map<String, List<SortSourceClusterInfo>> allTag2ClusterInfos =
clusterInfos.stream()
- .filter(dto -> dto.getClusterTags() != null)
- .filter(SortSourceClusterInfo::isConsumable)
- .filter(cluster ->
SUPPORTED_MQ_TYPE.contains(cluster.getType().toUpperCase(Locale.ROOT)))
-
.collect(Collectors.groupingBy(SortSourceClusterInfo::getClusterTags));
+ // reload all groups
+ groupInfos = configLoader.loadAllGroup()
+ .stream()
+ .collect(Collectors.toMap(SortSourceGroupInfo::getGroupId,
info -> info));
+
+ // reload all back up cluster
+ backupClusterTag =
configLoader.loadGroupBackupInfo(ClusterSwitch.BACKUP_CLUSTER_TAG)
+ .stream()
+
.collect(Collectors.toMap(InlongGroupExtEntity::getInlongGroupId,
InlongGroupExtEntity::getKeyValue));
+
+ // reload all back up group mq resource
+ backupGroupMqResource =
configLoader.loadGroupBackupInfo(ClusterSwitch.BACKUP_MQ_RESOURCE)
+ .stream()
+
.collect(Collectors.toMap(InlongGroupExtEntity::getInlongGroupId,
InlongGroupExtEntity::getKeyValue));
+
+ // reload all streams
+ allStreams = configLoader.loadAllStreams()
+ .stream()
+
.collect(Collectors.toMap(SortSourceStreamInfo::getInlongGroupId, stream ->
stream));
+
+ // reload all back up stream mq resource
+ backupStreamMqResource =
configLoader.loadStreamBackupInfo(ClusterSwitch.BACKUP_MQ_RESOURCE)
+ .stream()
+
.collect(Collectors.toMap(InlongStreamExtEntity::getInlongGroupId,
InlongStreamExtEntity::getKeyValue));
+ }
- // group clusters by name.
- Map<String, SortSourceClusterInfo> name2ClusterInfos =
clusterInfos.stream()
- .collect(Collectors.toMap(SortSourceClusterInfo::getName, info
-> info, (g1, g2) -> g1));
+ private void parseAll() {
// Prepare CacheZones for each cluster and task
Map<String, Map<String, String>> newMd5Map = new ConcurrentHashMap<>();
Map<String, Map<String, CacheZoneConfig>> newConfigMap = new
ConcurrentHashMap<>();
- groupMap.forEach((clusterName, task2Group) -> {
-
- // if there is no matched cluster name, just skip
- if (!name2ClusterInfos.containsKey(clusterName)) {
- return;
- }
- // find valid mq cluster list
- String clusterTag =
name2ClusterInfos.get(clusterName).getClusterTags();
- final Map<String, List<SortSourceClusterInfo>> validClusterInfos =
new ConcurrentHashMap<>();
- if (allTag2ClusterInfos.containsKey(clusterTag)) {
- validClusterInfos.put(clusterTag,
allTag2ClusterInfos.get(clusterTag));
- } else {
- validClusterInfos.putAll(allTag2ClusterInfos);
- }
+ groupMap.forEach((sortClusterName, task2GroupList) -> {
// prepare the new config and md5
Map<String, CacheZoneConfig> task2Config = new
ConcurrentHashMap<>();
Map<String, String> task2Md5 = new ConcurrentHashMap<>();
- task2Group.forEach((task, groupList) -> {
- // get topic properties under this cluster and task, group
them by group id.
- Map<String, Map<String, String>> group2topicProp = new
HashMap<>();
- allStreamInfos.stream().filter(stream ->
stream.getSortTaskName().equals(task)
- &&
stream.getSortClusterName().equals(clusterName)).forEach(
- sortSourceStreamInfo ->
group2topicProp.put(sortSourceStreamInfo.getGroupId(),
- sortSourceStreamInfo.getExtParamsMap()));
-
- Map<String, CacheZone> cacheZones;
+ task2GroupList.forEach((taskName, groupList) -> {
try {
- cacheZones = this.getCacheZones(groupList,
allId2GroupInfos, validClusterInfos, group2topicProp);
+ CacheZoneConfig cacheZoneConfig =
+ CacheZoneConfig.builder()
+ .sortClusterName(sortClusterName)
+ .sortTaskId(taskName)
+ .build();
+ Map<String, CacheZone> cacheZoneMap =
+ this.parseCacheZones(sortClusterName, taskName,
groupList);
+ cacheZoneConfig.setCacheZones(cacheZoneMap);
+
+ // prepare md5
+ String jsonStr = GSON.toJson(cacheZoneConfig);
+ String md5 = DigestUtils.md5Hex(jsonStr);
+ task2Config.put(taskName, cacheZoneConfig);
+ task2Md5.put(taskName, md5);
} catch (Throwable t) {
- LOGGER.error("fail to get cacheZones of clusterName {},
task {}", clusterName, task);
- return;
+ LOGGER.error("failed to parse sort source config of
sortCluster={}, task={}",
+ sortClusterName, taskName, t);
}
- CacheZoneConfig config = CacheZoneConfig.builder()
- .cacheZones(cacheZones)
- .sortClusterName(clusterName)
- .sortTaskId(task)
- .build();
- String jsonStr = GSON.toJson(config);
- String md5 = DigestUtils.md5Hex(jsonStr);
- task2Config.put(task, config);
- task2Md5.put(task, md5);
});
+ newConfigMap.put(sortClusterName, task2Config);
+ newMd5Map.put(sortClusterName, task2Md5);
- newConfigMap.put(clusterName, task2Config);
- newMd5Map.put(clusterName, task2Md5);
});
-
sortSourceConfigMap = newConfigMap;
sortSourceMd5Map = newMd5Map;
}
- private Map<String, CacheZone> getCacheZones(
- List<String> groupIdList,
- Map<String, SortSourceGroupInfo> allId2GroupInfos,
- Map<String, List<SortSourceClusterInfo>> allTag2ClusterInfos,
- Map<String, Map<String, String>> group2topicProp) {
+ private Map<String, CacheZone> parseCacheZones(
+ String sortClusterName,
+ String taskName,
+ List<String> groupIdList) {
- // stream of group info if group id exists.
- List<SortSourceGroupInfo> groupInfoStream = groupIdList.stream()
- .filter(allId2GroupInfos::containsKey)
- .map(allId2GroupInfos::get)
+ // get group infos
+ List<SortSourceGroupInfo> groupInfoList = groupIdList.stream()
+ .filter(groupInfos::containsKey)
+ .map(groupInfos::get)
.collect(Collectors.toList());
- // Group them by cluster tag.
- Map<String, List<SortSourceGroupInfo>> tag2GroupInfos =
groupInfoStream.stream()
+ // group them by cluster tag.
+ Map<String, List<SortSourceGroupInfo>> tag2GroupInfos =
groupInfoList.stream()
.collect(Collectors.groupingBy(SortSourceGroupInfo::getClusterTag));
- // Group them by back up cluster tag if both 2nd tag and 2nd topic
exist.
- Map<String, List<SortSourceGroupInfo>> backupTag2GroupInfos =
groupInfoStream.stream()
- .filter(group -> group.getBackupClusterTag() != null &&
group.getBackupTopic() != null)
-
.collect(Collectors.groupingBy(SortSourceGroupInfo::getBackupClusterTag));
+ // group them by second cluster tag.
+ Map<String, List<SortSourceGroupInfo>> backupTag2GroupInfos =
groupInfoList.stream()
+ .filter(info ->
backupClusterTag.containsKey(info.getGroupId()))
+ .collect(Collectors.groupingBy(info ->
backupClusterTag.get(info.getGroupId())));
- // get cache zone list.
- List<CacheZone> firstTagCacheZoneList =
- this.getCacheZoneListByTag(tag2GroupInfos,
allTag2ClusterInfos, group2topicProp, false);
- List<CacheZone> backupTagCacheZoneList =
- this.getCacheZoneListByTag(backupTag2GroupInfos,
allTag2ClusterInfos, group2topicProp, true);
+ List<CacheZone> cacheZones = this.parseCacheZonesByTag(tag2GroupInfos,
false);
+ List<CacheZone> backupCacheZones =
this.parseCacheZonesByTag(backupTag2GroupInfos, true);
- // combine two cache zone list, and group by cache zone name.
- return Stream.of(firstTagCacheZoneList, backupTagCacheZoneList)
+ return Stream.of(cacheZones, backupCacheZones)
.flatMap(Collection::stream)
.collect(Collectors.toMap(
CacheZone::getZoneName,
@@ -308,86 +305,63 @@ public class SortSourceServiceImpl implements
SortSourceService {
);
}
- private List<CacheZone> getCacheZoneListByTag(
- Map<String, List<SortSourceGroupInfo>> tag2GroupInfos,
- Map<String, List<SortSourceClusterInfo>> allTag2ClusterInfos,
- Map<String, Map<String, String>> group2topicProp,
- boolean isBackupTag) {
-
- // Tags of groups
- List<String> tags = new ArrayList<>(tag2GroupInfos.keySet());
+ private List<CacheZone> parseCacheZonesByTag(Map<String,
List<SortSourceGroupInfo>> tag2Groups, boolean isBackup) {
- // Clusters that related to these tags
- Map<String, List<SortSourceClusterInfo>> tag2ClusterInfos = new
HashMap<>();
- allTag2ClusterInfos.entrySet().stream().filter(entry ->
tag2GroupInfos.containsKey(entry.getKey()))
- .forEach(entry -> tag2ClusterInfos.put(entry.getKey(),
entry.getValue()));
-
- // get CacheZone list
- return tags.stream()
- .filter(tag2ClusterInfos::containsKey)
+ return tag2Groups.keySet().stream()
+ .filter(mqClusters::containsKey)
.flatMap(tag -> {
- List<SortSourceGroupInfo> groups = tag2GroupInfos.get(tag);
- List<SortSourceClusterInfo> clusters =
tag2ClusterInfos.get(tag);
+ List<SortSourceGroupInfo> groups = tag2Groups.get(tag);
+ List<SortSourceClusterInfo> clusters = mqClusters.get(tag);
return clusters.stream()
.map(cluster -> {
CacheZone zone = null;
try {
- zone = this.getCacheZone(groups, cluster,
group2topicProp, isBackupTag);
+ zone = this.parseCacheZone(groups,
cluster, isBackup);
} catch (IllegalStateException e) {
LOGGER.error("fail to init cache zone for
cluster " + cluster, e);
}
return zone;
});
})
- .filter(Objects::nonNull)
.collect(Collectors.toList());
}
- private CacheZone getCacheZone(
+ private CacheZone parseCacheZone(
List<SortSourceGroupInfo> groups,
SortSourceClusterInfo cluster,
- Map<String, Map<String, String>> group2topicProp,
boolean isBackupTag) {
+ switch (cluster.getType()) {
+ case ClusterType.PULSAR: return parsePulsarZone(groups, cluster,
isBackupTag);
+ default:
+ throw new BusinessException(String.format("do not support
cluster type=%s of cluster=%s",
+ cluster.getType(), cluster));
+ }
+ }
- // get basic Cache zone fields
+ private CacheZone parsePulsarZone(
+ List<SortSourceGroupInfo> groups,
+ SortSourceClusterInfo cluster,
+ boolean isBackupTag) {
Map<String, String> param = cluster.getExtParamsMap();
- String serviceUrl = Optional.ofNullable(param.get(KEY_SERVICE_URL))
- .orElseThrow(
- () -> new IllegalStateException(("there is no
serviceUrl for cluster " + cluster.getName())));
String tenant = param.get(KEY_TENANT);
- String namespace = param.get(KEY_NAME_SPACE);
- String authentication =
Optional.ofNullable(param.get(KEY_AUTH)).orElse("");
-
- List<Topic> topics = groups.stream()
- .map(groupInfo -> getTopic(groupInfo, tenant, namespace,
- group2topicProp.get(groupInfo.getGroupId()),
isBackupTag))
+ String auth = param.get(KEY_AUTH);
+ List<Topic> sdkTopics = groups.stream()
+ .map(info -> {
+ String namespace = info.getMqResource();
+ String topic =
allStreams.get(info.getGroupId()).getMqResource();
+ if (isBackupTag) {
+ namespace =
Optional.ofNullable(backupGroupMqResource.get(info.getGroupId())).orElse(namespace);
+ topic =
Optional.ofNullable(backupStreamMqResource.get(info.getGroupId())).orElse(topic);
+ }
+ String fullTopic =
tenant.concat("/").concat(namespace).concat("/").concat(topic);
+ return Topic.builder().topic(fullTopic).build();
+ })
.collect(Collectors.toList());
-
return CacheZone.builder()
- .serviceUrl(serviceUrl)
- .authentication(authentication)
- .cacheZoneProperties(param)
.zoneName(cluster.getName())
- .zoneType(cluster.getType())
- .topics(topics)
- .build();
- }
-
- private Topic getTopic(
- SortSourceGroupInfo groupInfo,
- String tenant,
- String namespace,
- Map<String, String> topicProperties,
- boolean isBackupTag) {
-
- String topic = isBackupTag ? groupInfo.getBackupTopic() :
groupInfo.getTopic();
- StringBuilder fullTopic = new StringBuilder();
- Optional.ofNullable(tenant).ifPresent(t ->
fullTopic.append(t).append("/"));
- Optional.ofNullable(namespace).ifPresent(n ->
fullTopic.append(n).append("/"));
- fullTopic.append(topic);
- return Topic.builder()
- .topic(fullTopic.toString())
- .topicProperties(topicProperties)
+ .serviceUrl(cluster.getUrl())
+ .topics(sdkTopics)
+ .authentication(auth)
.build();
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
index dc09f2be3..908e61220 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupOperator4Pulsar.java
@@ -135,16 +135,17 @@ public class InlongGroupOperator4Pulsar extends
AbstractGroupOperator {
// set backup topics, each inlong stream is associated with a Pulsar
topic
List<InlongStreamBriefInfo> streamTopics =
streamService.getTopicList(groupId);
- streamTopics.forEach(stream -> {
- InlongStreamExtEntity streamExtEntity =
streamExtMapper.selectByKey(groupId, stream.getInlongStreamId(),
- BACKUP_MQ_RESOURCE);
- if (streamExtEntity != null &&
StringUtils.isNotBlank(streamExtEntity.getKeyValue())) {
- stream.setMqResource(streamExtEntity.getKeyValue());
- }
- });
List<String> topics = streamTopics.stream()
- .map(InlongStreamBriefInfo::getMqResource)
+ .map(stream -> {
+ InlongStreamExtEntity streamExtEntity =
streamExtMapper.selectByKey(groupId,
+ stream.getInlongStreamId(), BACKUP_MQ_RESOURCE);
+ if (streamExtEntity != null &&
StringUtils.isNotBlank(streamExtEntity.getKeyValue())) {
+ return streamExtEntity.getKeyValue();
+ }
+ return stream.getMqResource();
+ })
.collect(Collectors.toList());
+
topicInfo.setTopics(topics);
return topicInfo;
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
index ddcb6f0c5..f89fd0f70 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/group/InlongGroupServiceImpl.java
@@ -389,7 +389,7 @@ public class InlongGroupServiceImpl implements
InlongGroupService {
public InlongGroupTopicInfo getBackupTopic(String groupId) {
// backup topic info saved in the ext table
InlongGroupExtEntity extEntity =
groupExtMapper.selectByUniqueKey(groupId, BACKUP_CLUSTER_TAG);
- if (StringUtils.isBlank(extEntity.getKeyValue())) {
+ if (extEntity == null || StringUtils.isBlank(extEntity.getKeyValue()))
{
LOGGER.warn("not found any backup topic for groupId={}", groupId);
return null;
}
diff --git
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
index b65735236..f8eb9eff8 100644
---
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
+++
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/sort/SortServiceImplTest.java
@@ -17,19 +17,28 @@
package org.apache.inlong.manager.service.sort;
+import org.apache.inlong.common.constant.ClusterSwitch;
import org.apache.inlong.common.pojo.sdk.SortSourceConfigResponse;
import org.apache.inlong.common.pojo.sortstandalone.SortClusterResponse;
import org.apache.inlong.manager.common.consts.InlongConstants;
+import org.apache.inlong.manager.common.enums.ClusterType;
import org.apache.inlong.manager.dao.entity.DataNodeEntity;
import org.apache.inlong.manager.dao.entity.InlongClusterEntity;
-import org.apache.inlong.manager.dao.entity.InlongGroupEntity;
import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
import org.apache.inlong.manager.dao.mapper.DataNodeEntityMapper;
import org.apache.inlong.manager.dao.mapper.InlongClusterEntityMapper;
import org.apache.inlong.manager.dao.mapper.InlongGroupEntityMapper;
import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
+import org.apache.inlong.manager.pojo.cluster.pulsar.PulsarClusterRequest;
+import org.apache.inlong.manager.pojo.group.InlongGroupExtInfo;
+import org.apache.inlong.manager.pojo.group.pulsar.InlongPulsarRequest;
+import org.apache.inlong.manager.pojo.stream.InlongStreamExtInfo;
+import org.apache.inlong.manager.pojo.stream.InlongStreamRequest;
import org.apache.inlong.manager.service.ServiceBaseTest;
+import org.apache.inlong.manager.service.cluster.InlongClusterService;
import org.apache.inlong.manager.service.core.SortService;
+import org.apache.inlong.manager.service.group.InlongGroupService;
+import org.apache.inlong.manager.service.stream.InlongStreamService;
import org.json.JSONObject;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeEach;
@@ -40,7 +49,9 @@ import org.junit.jupiter.api.TestMethodOrder;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.transaction.annotation.Transactional;
+import java.util.ArrayList;
import java.util.Date;
+import java.util.List;
/**
* Sort service test for {@link SortService}
@@ -51,7 +62,9 @@ public class SortServiceImplTest extends ServiceBaseTest {
private static final String TEST_CLUSTER = "testCluster";
private static final String TEST_TASK = "testTask";
private static final String TEST_GROUP = "testGroup";
+ private static final String TEST_STREAM = "1";
private static final String TEST_TAG = "testTag";
+ private static final String BACK_UP_TAG = "testBackupTag";
private static final String TEST_TOPIC = "testTopic";
private static final String TEST_SINK_TYPE = "testSinkType";
private static final String TEST_CREATOR = "testUser";
@@ -66,6 +79,12 @@ public class SortServiceImplTest extends ServiceBaseTest {
private DataNodeEntityMapper dataNodeEntityMapper;
@Autowired
private SortService sortService;
+ @Autowired
+ private InlongClusterService clusterService;
+ @Autowired
+ private InlongGroupService groupService;
+ @Autowired
+ private InlongStreamService streamService;
@Test
@Order(1)
@@ -173,10 +192,12 @@ public class SortServiceImplTest extends ServiceBaseTest {
@BeforeEach
private void prepareAll() {
this.prepareCluster(TEST_CLUSTER);
- this.prepareTask(TEST_TASK, TEST_GROUP, TEST_CLUSTER);
- this.prepareGroupId(TEST_GROUP);
- this.preparePulsar("testPulsar", true);
+ this.preparePulsar("testPulsar", true, TEST_TAG);
+ this.preparePulsar("testPulsar2", true, BACK_UP_TAG);
this.prepareDataNode(TEST_TASK);
+ this.prepareGroupId(TEST_GROUP);
+ this.prepareStreamId(TEST_GROUP, TEST_STREAM);
+ this.prepareTask(TEST_TASK, TEST_GROUP, TEST_CLUSTER);
}
private void prepareDataNode(String taskName) {
@@ -195,20 +216,50 @@ public class SortServiceImplTest extends ServiceBaseTest {
}
private void prepareGroupId(String groupId) {
- InlongGroupEntity entity = new InlongGroupEntity();
- entity.setInlongGroupId(groupId);
- entity.setInlongClusterTag(TEST_TAG);
- entity.setMqResource(TEST_TOPIC);
- entity.setMqType("PULSAR");
- entity.setName("testName");
- entity.setDescription("testDescription");
- entity.setCreator(TEST_CREATOR);
- entity.setInCharges(TEST_CREATOR);
- entity.setCreateTime(new Date());
- entity.setModifyTime(new Date());
- entity.setIsDeleted(InlongConstants.UN_DELETED);
- entity.setVersion(InlongConstants.INITIAL_VERSION);
- inlongGroupEntityMapper.insert(entity);
+ InlongPulsarRequest request = new InlongPulsarRequest();
+ request.setInlongGroupId(groupId);
+ request.setMqResource("test_namespace");
+ request.setInlongClusterTag(TEST_TAG);
+ request.setVersion(InlongConstants.INITIAL_VERSION);
+ request.setName("test_group_name");
+ request.setMqType(ClusterType.PULSAR);
+ request.setInCharges(TEST_CREATOR);
+ List<InlongGroupExtInfo> extList = new ArrayList<>();
+ InlongGroupExtInfo ext1 = InlongGroupExtInfo
+ .builder()
+ .inlongGroupId(groupId)
+ .keyName(ClusterSwitch.BACKUP_CLUSTER_TAG)
+ .keyValue(BACK_UP_TAG)
+ .build();
+ InlongGroupExtInfo ext2 = InlongGroupExtInfo
+ .builder()
+ .inlongGroupId(groupId)
+ .keyName(ClusterSwitch.BACKUP_MQ_RESOURCE)
+ .keyValue("backup_name")
+ .build();
+
+ extList.add(ext1);
+ extList.add(ext2);
+ request.setExtList(extList);
+ groupService.save(request, "test operator");
+ }
+
+ private void prepareStreamId(String groupId, String streamId) {
+ InlongStreamRequest request = new InlongStreamRequest();
+ request.setInlongGroupId(groupId);
+ request.setInlongStreamId(streamId);
+ request.setName("test_stream_name");
+ request.setMqResource(TEST_TOPIC);
+ request.setVersion(InlongConstants.INITIAL_VERSION);
+ List<InlongStreamExtInfo> extInfos = new ArrayList<>();
+ InlongStreamExtInfo ext = new InlongStreamExtInfo();
+ extInfos.add(ext);
+ ext.setInlongStreamId(streamId);
+ ext.setInlongGroupId(groupId);
+ ext.setKeyName(ClusterSwitch.BACKUP_MQ_RESOURCE);
+ ext.setKeyValue("backup_topic");
+ request.setExtList(extInfos);
+ streamService.save(request, "test_operator");
}
private void prepareCluster(String clusterName) {
@@ -226,25 +277,21 @@ public class SortServiceImplTest extends ServiceBaseTest {
clusterEntityMapper.insert(entity);
}
- private void preparePulsar(String pulsarName, boolean isConsumable) {
- InlongClusterEntity entity = new InlongClusterEntity();
- entity.setName(pulsarName);
- entity.setType("PULSAR");
- entity.setClusterTags(TEST_TAG);
- entity.setCreator(TEST_CREATOR);
- entity.setInCharges(TEST_CREATOR);
- Date now = new Date();
- entity.setCreateTime(now);
- entity.setModifyTime(now);
- entity.setIsDeleted(InlongConstants.UN_DELETED);
- entity.setVersion(InlongConstants.INITIAL_VERSION);
+ private void preparePulsar(String pulsarName, boolean isConsumable, String
tag) {
+ PulsarClusterRequest request = new PulsarClusterRequest();
+ request.setUrl("testServiceUrl");
+ request.setName(pulsarName);
+ request.setType(ClusterType.PULSAR);
+ request.setClusterTags(tag);
+ request.setInCharges(TEST_CREATOR);
+ request.setVersion(InlongConstants.INITIAL_VERSION);
String extTag = "zone=" + TEST_TAG
+ "&producer=true"
+ "&consumer=" + (isConsumable ? "true" : "false");
- entity.setExtTag(extTag);
-
entity.setExtParams("{\"tenant\":\"testTenant\",\"namespace\":\"testNS\",\"serviceUrl\":\"testServiceUrl\","
+ request.setExtTag(extTag);
+ request.setExtParams("{\"tenant\":\"testTenant\","
+
"\"authentication\":\"testAuth\",\"adminUrl\":\"testAdmin\"}");
- clusterEntityMapper.insert(entity);
+ clusterService.save(request, "test operator");
}
private void prepareTask(String taskName, String groupId, String
clusterName) {