This is an automated email from the ASF dual-hosted git repository.
healchow pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-inlong.git
The following commit(s) were added to refs/heads/master by this push:
new d1ff90ba9 [INLONG-4223][Manager] Refactor the consumption table
structure (#4226)
d1ff90ba9 is described below
commit d1ff90ba99cd81c8a426916a72cf6ef5a657f542
Author: healchow <[email protected]>
AuthorDate: Tue May 17 11:55:57 2022 +0800
[INLONG-4223][Manager] Refactor the consumption table structure (#4226)
---
.../inlong/manager/common/enums/ErrorCodeEnum.java | 6 +-
.../common/pojo/consumption/ConsumptionInfo.java | 24 ++--
.../common/pojo/consumption/ConsumptionListVo.java | 14 +-
.../pojo/consumption/ConsumptionMqExtBase.java | 6 +-
.../pojo/consumption/ConsumptionPulsarInfo.java | 2 +-
.../common/pojo/consumption/ConsumptionQuery.java | 21 +--
.../pojo/workflow/form/ConsumptionApproveForm.java | 6 +-
.../workflow/form/NewConsumptionProcessForm.java | 2 +-
.../manager/dao/entity/ConsumptionEntity.java | 8 +-
.../dao/entity/ConsumptionPulsarEntity.java | 5 +-
.../main/resources/mappers/ClusterSetMapper.xml | 14 +-
.../resources/mappers/ConsumptionEntityMapper.xml | 151 ++++++++-------------
.../mappers/ConsumptionPulsarEntityMapper.xml | 24 ++--
.../manager/service/core/ConsumptionService.java | 6 +-
.../service/core/impl/ConsumptionServiceImpl.java | 31 ++---
.../manager/service/core/mq/MiddlewareFactory.java | 3 +-
.../ConsumptionCompleteProcessListener.java | 18 +--
.../listener/ConsumptionPassTaskListener.java | 16 +--
.../service/core/impl/ConsumptionServiceTest.java | 6 +-
.../core/impl/InlongGroupProcessOperationTest.java | 3 +-
.../main/resources/sql/apache_inlong_manager.sql | 48 ++++---
.../manager-web/sql/apache_inlong_manager.sql | 48 ++++---
22 files changed, 197 insertions(+), 265 deletions(-)
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/ErrorCodeEnum.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/ErrorCodeEnum.java
index 8a6a8be3a..d6eace625 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/ErrorCodeEnum.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/enums/ErrorCodeEnum.java
@@ -40,7 +40,7 @@ public enum ErrorCodeEnum {
GROUP_DELETE_NOT_ALLOWED(1007, "The current inlong group status does not
support deletion"),
GROUP_ID_UPDATE_NOT_ALLOWED(1008, "The current inlong group status does
not support modifying the group id"),
GROUP_MIDDLEWARE_UPDATE_NOT_ALLOWED(1011,
- "The current inlong group status does not support modifying the
middleware type"),
+ "The current inlong group status does not support modifying the MQ
type"),
GROUP_NAME_UPDATE_NOT_ALLOWED(1012, "The current inlong group status does
not support modifying the name"),
GROUP_INFO_INCONSISTENT(1013, "The inlong group info is inconsistent,
please contact the administrator"),
GROUP_MODE_UNSUPPORTED(1014, "The current inlong group mode only support
light, normal"),
@@ -48,7 +48,7 @@ public enum ErrorCodeEnum {
OPT_NOT_ALLOWED_BY_STATUS(1021,
"The current inlong group status does not allow
adding/modifying/deleting related info"),
- MQ_TYPE_NOT_SUPPORTED(1021, "MIDDLEWARE_TYPE_NOT_SUPPORTED"),
+ MQ_TYPE_NOT_SUPPORTED(1021, "MQ_TYPE_NOT_SUPPORTED"),
CLUSTER_NOT_FOUND(1101, "Cluster information does not exist"),
@@ -100,7 +100,7 @@ public enum ErrorCodeEnum {
WORKFLOW_EXE_FAILED(4000, "Workflow execution exception"),
- CONSUMER_GROUP_NAME_DUPLICATED(2600, "The consumer group already exists in
the cluster"),
+ CONSUMER_GROUP_DUPLICATED(2600, "The consumer group already exists"),
CONSUMER_GROUP_CREATE_FAILED(2601, "Failed to create tube consumer group"),
TUBE_GROUP_CREATE_FAILED(2602, "Create Tube consumer group failed"),
PULSAR_GROUP_CREATE_FAILED(2603, "Create Pulsar consumer group failed"),
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionInfo.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionInfo.java
index 397789e85..ecef9c7b0 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionInfo.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionInfo.java
@@ -20,16 +20,17 @@ package org.apache.inlong.manager.common.pojo.consumption;
import com.fasterxml.jackson.annotation.JsonIgnore;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
-import java.util.Date;
-import javax.validation.constraints.AssertTrue;
-import javax.validation.constraints.NotBlank;
-import javax.validation.constraints.NotNull;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.apache.commons.lang3.StringUtils;
+import javax.validation.constraints.AssertTrue;
+import javax.validation.constraints.NotBlank;
+import javax.validation.constraints.NotNull;
+import java.util.Date;
+
/**
* Data consumption info
*/
@@ -43,13 +44,9 @@ public class ConsumptionInfo {
@ApiModelProperty(value = "key id")
private Integer id;
- @ApiModelProperty(value = "consumer group id")
- @NotBlank(message = "consumerGroupId cannot be null")
- private String consumerGroupId;
-
- @ApiModelProperty(value = "consumer group name: only support [a-zA-Z0-9_]")
- @NotBlank(message = "consumerGroupName cannot be null")
- private String consumerGroupName;
+ @ApiModelProperty(value = "consumer group: only support [a-zA-Z0-9_]")
+ @NotBlank(message = "consumerGroup cannot be null")
+ private String consumerGroup;
@ApiModelProperty(value = "consumption in charge")
@NotNull(message = "inCharges cannot be null")
@@ -59,9 +56,8 @@ public class ConsumptionInfo {
@NotBlank(message = "inlong group id cannot be null")
private String inlongGroupId;
- @ApiModelProperty(value = "Middleware type, high throughput: TUBE, high
consistency: PULSAR")
- // @NotBlank(message = "middlewareType cannot be null")
- private String middlewareType;
+ @ApiModelProperty(value = "MQ type, high throughput: TUBE, high
consistency: PULSAR")
+ private String mqType;
@ApiModelProperty(value = "consumption target topic")
private String topic;
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionListVo.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionListVo.java
index ab634cdf4..f789f119f 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionListVo.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionListVo.java
@@ -19,12 +19,13 @@ package org.apache.inlong.manager.common.pojo.consumption;
import io.swagger.annotations.ApiModel;
import io.swagger.annotations.ApiModelProperty;
-import java.util.Date;
import lombok.AllArgsConstructor;
import lombok.Builder;
import lombok.Data;
import lombok.NoArgsConstructor;
+import java.util.Date;
+
/**
* Data consumption list
*/
@@ -38,11 +39,8 @@ public class ConsumptionListVo {
@ApiModelProperty(value = "Primary key")
private Integer id;
- @ApiModelProperty(value = "Consumer group name-lowercase letters, numbers,
underscores")
- private String consumerGroupName;
-
- @ApiModelProperty(value = "Consumer Group ID")
- private String consumerGroupId;
+ @ApiModelProperty(value = "Consumer Group")
+ private String consumerGroup;
@ApiModelProperty(value = "Person in charge of consumption")
private String inCharges;
@@ -50,8 +48,8 @@ public class ConsumptionListVo {
@ApiModelProperty(value = "Consumption target inlong group id")
private String inlongGroupId;
- @ApiModelProperty(value = "Middleware type, high throughput: TUBE, high
consistency: PULSAR")
- private String middlewareType;
+ @ApiModelProperty(value = "MQ type, high throughput: TUBE, high
consistency: PULSAR")
+ private String mqType;
@ApiModelProperty(value = "Consumption target TOPIC")
private String topic;
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionMqExtBase.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionMqExtBase.java
index a49a1e22b..6bbb13e00 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionMqExtBase.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionMqExtBase.java
@@ -29,7 +29,7 @@ import lombok.Data;
*/
@Data
@ApiModel("Extended consumption information of different MQs")
-@JsonTypeInfo(use = Id.NAME, visible = true, property = "middlewareType",
defaultImpl = ConsumptionMqExtBase.class)
+@JsonTypeInfo(use = Id.NAME, visible = true, property = "mqType", defaultImpl
= ConsumptionMqExtBase.class)
@JsonSubTypes({
@JsonSubTypes.Type(value = ConsumptionPulsarInfo.class, name =
"PULSAR"),
@JsonSubTypes.Type(value = ConsumptionPulsarInfo.class, name =
"TDMQ_PULSAR")
@@ -51,6 +51,6 @@ public class ConsumptionMqExtBase {
@ApiModelProperty("Whether to delete, 0: not deleted, 1: deleted")
private Integer isDeleted = 0;
- @ApiModelProperty("The middleware type of MQ")
- private String middlewareType;
+ @ApiModelProperty("The type of MQ")
+ private String mqType;
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionPulsarInfo.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionPulsarInfo.java
index 3c49ba609..d102ea4fd 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionPulsarInfo.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionPulsarInfo.java
@@ -34,7 +34,7 @@ import org.apache.inlong.manager.common.enums.MQType;
public class ConsumptionPulsarInfo extends ConsumptionMqExtBase {
public ConsumptionPulsarInfo() {
- this.setMiddlewareType(MQType.PULSAR.getType());
+ this.setMqType(MQType.PULSAR.getType());
}
@ApiModelProperty("Whether to configure the dead letter queue, 0: do not
configure, 1: configure")
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionQuery.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionQuery.java
index cac432fc5..8d53f65a6 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionQuery.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/consumption/ConsumptionQuery.java
@@ -31,14 +31,8 @@ import org.apache.inlong.manager.common.beans.PageRequest;
@ApiModel("Data consumption query conditions")
public class ConsumptionQuery extends PageRequest {
- @ApiModelProperty(value = "Consumer Group Name")
- private String consumerGroupName;
-
- @ApiModelProperty(value = "Consumer Group ID")
- private String consumerGroupId;
-
- @ApiModelProperty(value = "Consumer Group Id is fuzzy")
- private String consumerGroupIdLike;
+ @ApiModelProperty(value = "Consumer Group")
+ private String consumerGroup;
@ApiModelProperty(value = "Person in charge of consumption")
private String inCharges;
@@ -46,15 +40,12 @@ public class ConsumptionQuery extends PageRequest {
@ApiModelProperty(value = "Consumption target inlong group id")
private String inlongGroupId;
- @ApiModelProperty(value = "Middleware type, high throughput: TUBE, high
consistency: PULSAR")
- private String middlewareType;
+ @ApiModelProperty(value = "MQ type, high throughput: TUBE, high
consistency: PULSAR")
+ private String mqType;
- @ApiModelProperty(value = "Consumption target TOPIC")
+ @ApiModelProperty(value = "Consumption target Topic")
private String topic;
- @ApiModelProperty(value = "Fuzzy matching of consumption target TOPIC")
- private String topicLike;
-
@ApiModelProperty(value = "Whether to filter consumption")
private Boolean filterEnabled;
@@ -78,6 +69,6 @@ public class ConsumptionQuery extends PageRequest {
@ApiModelProperty(value = "Weather current user have admin role", hidden =
true)
private Boolean isAdminRole;
- @ApiModelProperty(value = "Fuzzy query keyword, fuzzy query topic,
consumer group ID")
+ @ApiModelProperty(value = "Fuzzy query keyword, topic, or consumer group")
private String keyword;
}
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/ConsumptionApproveForm.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/ConsumptionApproveForm.java
index da370dc1e..839316871 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/ConsumptionApproveForm.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/ConsumptionApproveForm.java
@@ -32,13 +32,13 @@ public class ConsumptionApproveForm extends BaseTaskForm {
public static final String FORM_NAME = "ConsumptionApproveForm";
- @ApiModelProperty("Consumer group ID")
- private String consumerGroupId;
+ @ApiModelProperty("Consumer group")
+ private String consumerGroup;
@Override
public void validate() throws FormValidateException {
- Preconditions.checkNotEmpty(consumerGroupId, "Consumer group cannot be
empty");
+ Preconditions.checkNotEmpty(consumerGroup, "Consumer group cannot be
empty");
}
@Override
diff --git
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/NewConsumptionProcessForm.java
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/NewConsumptionProcessForm.java
index aeeddff91..f174de183 100644
---
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/NewConsumptionProcessForm.java
+++
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/workflow/form/NewConsumptionProcessForm.java
@@ -51,7 +51,7 @@ public class NewConsumptionProcessForm extends
BaseProcessForm {
@Override
public String getInlongGroupId() {
- return consumptionInfo.getConsumerGroupId();
+ return consumptionInfo.getConsumerGroup();
}
@Override
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionEntity.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionEntity.java
index 6b1ad14b4..9e8cda5d4 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionEntity.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionEntity.java
@@ -17,9 +17,10 @@
package org.apache.inlong.manager.dao.entity;
+import lombok.Data;
+
import java.io.Serializable;
import java.util.Date;
-import lombok.Data;
/**
* Data consumption table
@@ -29,11 +30,10 @@ public class ConsumptionEntity implements Serializable {
private static final long serialVersionUID = 1L;
private Integer id;
- private String consumerGroupName;
- private String consumerGroupId;
+ private String consumerGroup;
private String inCharges;
private String inlongGroupId;
- private String middlewareType;
+ private String mqType;
private String topic;
private Integer filterEnabled;
private String inlongStreamId;
diff --git
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionPulsarEntity.java
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionPulsarEntity.java
index a524d592a..32e9ed129 100644
---
a/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionPulsarEntity.java
+++
b/inlong-manager/manager-dao/src/main/java/org/apache/inlong/manager/dao/entity/ConsumptionPulsarEntity.java
@@ -17,9 +17,10 @@
package org.apache.inlong.manager.dao.entity;
-import java.io.Serializable;
import lombok.Data;
+import java.io.Serializable;
+
/**
* Data consumption table, including group name, group id, topic, etc
*/
@@ -29,7 +30,7 @@ public class ConsumptionPulsarEntity implements Serializable {
private static final long serialVersionUID = 1L;
private Integer id;
private Integer consumptionId;
- private String consumerGroupId;
+ private String consumerGroup;
private String inlongGroupId;
private Integer isDlq;
private String deadLetterTopic;
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/ClusterSetMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/ClusterSetMapper.xml
index 8cfeb6ef0..2eb1b176e 100644
--- a/inlong-manager/manager-dao/src/main/resources/mappers/ClusterSetMapper.xml
+++ b/inlong-manager/manager-dao/src/main/resources/mappers/ClusterSetMapper.xml
@@ -35,14 +35,14 @@
where is_deleted = 0
</select>
<select id="selectInlongId"
resultType="org.apache.inlong.manager.dao.entity.InLongId">
- select biz.inlong_stream_id as inlong_id,
- biz.mq_resource_obj as topic,
- concat('fieldDelimiter=',biz.data_separator) as params,
- c.set_name as set_name
- from inlong_stream biz,
+ select stream.inlong_stream_id as inlong_id,
+ stream.mq_resource as topic,
+ concat('fieldDelimiter=', stream.data_separator) as params,
+ c.set_name as set_name
+ from inlong_stream stream,
cluster_set_inlongid c
- where biz.is_deleted = 0
- and biz.inlong_group_id = c.inlong_group_id
+ where stream.is_deleted = 0
+ and stream.inlong_group_id = c.inlong_group_id
</select>
<select id="selectCacheCluster"
resultType="org.apache.inlong.manager.dao.entity.CacheCluster">
select cluster_name, set_name, zone
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionEntityMapper.xml
index e0fd2ada4..9b414a277 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionEntityMapper.xml
@@ -23,11 +23,10 @@
<mapper
namespace="org.apache.inlong.manager.dao.mapper.ConsumptionEntityMapper">
<resultMap id="BaseResultMap"
type="org.apache.inlong.manager.dao.entity.ConsumptionEntity">
<id column="id" jdbcType="INTEGER" property="id"/>
- <result column="consumer_group_name" jdbcType="VARCHAR"
property="consumerGroupName"/>
- <result column="consumer_group_id" jdbcType="VARCHAR"
property="consumerGroupId"/>
+ <result column="consumer_group" jdbcType="VARCHAR"
property="consumerGroup"/>
<result column="in_charges" jdbcType="VARCHAR" property="inCharges"/>
<result column="inlong_group_id" jdbcType="VARCHAR"
property="inlongGroupId"/>
- <result column="middleware_type" jdbcType="VARCHAR"
property="middlewareType"/>
+ <result column="mq_type" jdbcType="VARCHAR" property="mqType"/>
<result column="topic" jdbcType="VARCHAR" property="topic"/>
<result column="filter_enabled" jdbcType="INTEGER"
property="filterEnabled"/>
<result column="inlong_stream_id" jdbcType="VARCHAR"
property="inlongStreamId"/>
@@ -39,8 +38,8 @@
<result column="is_deleted" jdbcType="INTEGER" property="isDeleted"/>
</resultMap>
<sql id="Base_Column_List">
- id, consumer_group_name, consumer_group_id, in_charges,
inlong_group_id,
- middleware_type, topic, filter_enabled, inlong_stream_id,
+ id, consumer_group, in_charges, inlong_group_id,
+ mq_type, topic, filter_enabled, inlong_stream_id,
status, is_deleted, creator, modifier, create_time, modify_time
</sql>
<select id="selectByPrimaryKey" parameterType="java.lang.Integer"
@@ -56,44 +55,39 @@
from consumption
where inlong_group_id = #{groupId, jdbcType=VARCHAR}
and topic = #{topic, jdbcType=VARCHAR}
- and consumer_group_id = #{consumerGroup, jdbcType=VARCHAR}
+ and consumer_group = #{consumerGroup, jdbcType=VARCHAR}
and is_deleted = 0
</select>
<delete id="deleteByPrimaryKey" parameterType="java.lang.Integer">
update consumption
- set is_deleted=id
+ set is_deleted = id
where id = #{id,jdbcType=INTEGER}
</delete>
<insert id="insert" useGeneratedKeys="true" keyProperty="id"
parameterType="org.apache.inlong.manager.dao.entity.ConsumptionEntity">
- insert into consumption (id, consumer_group_name,
- consumer_group_id, in_charges,
- inlong_group_id, middleware_type, topic,
+ insert into consumption (id, consumer_group, in_charges,
+ inlong_group_id, mq_type, topic,
filter_enabled, inlong_stream_id,
status, is_deleted,
creator, modifier,
create_time, modify_time)
- values (#{id,jdbcType=INTEGER}, #{consumerGroupName,jdbcType=VARCHAR},
- #{consumerGroupId,jdbcType=VARCHAR},
#{inCharges,jdbcType=VARCHAR},
- #{inlongGroupId,jdbcType=VARCHAR},
#{middlewareType,jdbcType=VARCHAR}, #{topic,jdbcType=VARCHAR},
+ values (#{id,jdbcType=INTEGER}, #{consumerGroup,jdbcType=VARCHAR},
#{inCharges,jdbcType=VARCHAR},
+ #{inlongGroupId,jdbcType=VARCHAR}, #{mqType,jdbcType=VARCHAR},
#{topic,jdbcType=VARCHAR},
#{filterEnabled,jdbcType=INTEGER},
#{inlongStreamId,jdbcType=VARCHAR},
#{status,jdbcType=INTEGER}, #{isDeleted,jdbcType=INTEGER},
#{creator,jdbcType=VARCHAR}, #{modifier,jdbcType=VARCHAR},
#{createTime,jdbcType=TIMESTAMP},
#{modifyTime,jdbcType=TIMESTAMP})
</insert>
- <insert id="insertSelective"
+ <insert id="insertSelective" useGeneratedKeys="true" keyProperty="id"
parameterType="org.apache.inlong.manager.dao.entity.ConsumptionEntity">
insert into consumption
<trim prefix="(" suffix=")" suffixOverrides=",">
<if test="id != null">
id,
</if>
- <if test="consumerGroupName != null">
- consumer_group_name,
- </if>
- <if test="consumerGroupId != null">
- consumer_group_id,
+ <if test="consumerGroup != null">
+ consumer_group,
</if>
<if test="inCharges != null">
in_charges,
@@ -101,8 +95,8 @@
<if test="inlongGroupId != null">
inlong_group_id,
</if>
- <if test="middlewareType != null">
- middleware_type,
+ <if test="mqType != null">
+ mq_type,
</if>
<if test="topic != null">
topic,
@@ -136,11 +130,8 @@
<if test="id != null">
#{id,jdbcType=INTEGER},
</if>
- <if test="consumerGroupName != null">
- #{consumerGroupName,jdbcType=VARCHAR},
- </if>
- <if test="consumerGroupId != null">
- #{consumerGroupId,jdbcType=VARCHAR},
+ <if test="consumerGroup != null">
+ #{consumerGroup,jdbcType=VARCHAR},
</if>
<if test="inCharges != null">
#{inCharges,jdbcType=VARCHAR},
@@ -148,8 +139,8 @@
<if test="inlongGroupId != null">
#{inlongGroupId,jdbcType=VARCHAR},
</if>
- <if test="middlewareType != null">
- #{middlewareType,jdbcType=VARCHAR},
+ <if test="mqType != null">
+ #{mqType,jdbcType=VARCHAR},
</if>
<if test="topic != null">
#{topic,jdbcType=VARCHAR},
@@ -184,11 +175,8 @@
parameterType="org.apache.inlong.manager.dao.entity.ConsumptionEntity">
update consumption
<set>
- <if test="consumerGroupName != null">
- consumer_group_name = #{consumerGroupName,jdbcType=VARCHAR},
- </if>
- <if test="consumerGroupId != null">
- consumer_group_id = #{consumerGroupId,jdbcType=VARCHAR},
+ <if test="consumerGroup != null">
+ consumer_group = #{consumerGroup,jdbcType=VARCHAR},
</if>
<if test="inCharges != null">
in_charges = #{inCharges,jdbcType=VARCHAR},
@@ -196,8 +184,8 @@
<if test="inlongGroupId != null">
inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR},
</if>
- <if test="middlewareType != null">
- middleware_type = #{middlewareType,jdbcType=VARCHAR},
+ <if test="mqType != null">
+ mq_type = #{mqType,jdbcType=VARCHAR},
</if>
<if test="topic != null">
topic = #{topic,jdbcType=VARCHAR},
@@ -228,20 +216,19 @@
</update>
<update id="updateByPrimaryKey"
parameterType="org.apache.inlong.manager.dao.entity.ConsumptionEntity">
update consumption
- set consumer_group_name = #{consumerGroupName,jdbcType=VARCHAR},
- consumer_group_id = #{consumerGroupId,jdbcType=VARCHAR},
- in_charges = #{inCharges,jdbcType=VARCHAR},
- inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR},
- middleware_type = #{middlewareType,jdbcType=VARCHAR},
- topic = #{topic,jdbcType=VARCHAR},
- filter_enabled = #{filterEnabled,jdbcType=INTEGER},
- inlong_stream_id = #{inlongStreamId,jdbcType=VARCHAR},
- status = #{status,jdbcType=INTEGER},
- creator = #{creator,jdbcType=VARCHAR},
- modifier = #{modifier,jdbcType=VARCHAR},
- create_time = #{createTime,jdbcType=TIMESTAMP},
- modify_time = #{modifyTime,jdbcType=TIMESTAMP},
- is_deleted = #{isDeleted,jdbcType=INTEGER}
+ set consumer_group = #{consumerGroup,jdbcType=VARCHAR},
+ in_charges = #{inCharges,jdbcType=VARCHAR},
+ inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR},
+ mq_type = #{mqType,jdbcType=VARCHAR},
+ topic = #{topic,jdbcType=VARCHAR},
+ filter_enabled = #{filterEnabled,jdbcType=INTEGER},
+ inlong_stream_id = #{inlongStreamId,jdbcType=VARCHAR},
+ status = #{status,jdbcType=INTEGER},
+ creator = #{creator,jdbcType=VARCHAR},
+ modifier = #{modifier,jdbcType=VARCHAR},
+ create_time = #{createTime,jdbcType=TIMESTAMP},
+ modify_time = #{modifyTime,jdbcType=TIMESTAMP},
+ is_deleted = #{isDeleted,jdbcType=INTEGER}
where id = #{id,jdbcType=INTEGER}
and is_deleted = 0
</update>
@@ -253,28 +240,17 @@
c.*
from consumption c
where c.is_deleted=0
- <if test="consumerGroupName != null and consumerGroupName != ''">
- and c.consumer_group_name = #{consumerGroupName,jdbcType=VARCHAR}
- </if>
- <if test="consumerGroupId != null and consumerGroupId != ''">
- and c.consumer_group_id = #{consumerGroupId,jdbcType=VARCHAR}
- </if>
- <if test="consumerGroupIdLike != null and consumerGroupIdLike != ''">
- <bind name="consumerGroupIdLike" value="'%' +
_parameter.consumerGroupIdLike + '%'"/>
- and c.consumer_group_id like
#{consumerGroupIdLike,jdbcType=VARCHAR}
+ <if test="consumerGroup != null and consumerGroup != ''">
+ and c.consumer_group = #{consumerGroup,jdbcType=VARCHAR}
</if>
<if test="inlongGroupId != null and inlongGroupId != ''">
and c.inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR}
</if>
- <if test="middlewareType != null and middlewareType != ''">
- and c.middleware_type = #{middlewareType,jdbcType=VARCHAR}
+ <if test="mqType != null and mqType != ''">
+ and c.mq_type = #{mqType,jdbcType=VARCHAR}
</if>
<if test="topic != null and topic != ''">
- and c.topic = #{topic,jdbcType=VARCHAR}
- </if>
- <if test="topicLike != null and topicLike != ''">
- <bind name="topicLike" value="'%' + _parameter.topicLike + '%'"/>
- and c.topic like #{topicLike,jdbcType=VARCHAR}
+ and c.topic like CONCAT('%', #{topic}, '%')
</if>
<if test="filterEnabled != null">
and c.filter_enabled = #{filterEnabled,jdbcType=INTEGER}
@@ -294,22 +270,17 @@
<if test="isAdminRole == false">
and (
FIND_IN_SET(#{userName,jdbcType=VARCHAR},c.in_charges)
- or
- c.creator = #{userName,jdbcType=VARCHAR}
+ or c.creator = #{userName,jdbcType=VARCHAR}
)
</if>
<if test="keyword != null and keyword !=''">
- <bind name="keyword" value="'%' + _parameter.keyword + '%'"/>
- and (
- c.topic like #{keyword,jdbcType=VARCHAR}
- or c.consumer_group_id like #{keyword,jdbcType=VARCHAR}
- )
+ and (c.topic like CONCAT('%', #{keyword}, '%') or c.consumer_group
like CONCAT('%', #{keyword}, '%'))
</if>
- <if test="lastConsumptionStatus!=null and lastConsumptionStatus!=3">
- and cas.latest_record_state=#{lastConsumptionStatus,
jdbcType=INTEGER}
+ <if test="lastConsumptionStatus != null and lastConsumptionStatus !=
3">
+ and cas.latest_record_state = #{lastConsumptionStatus,
jdbcType=INTEGER}
</if>
- <if test="lastConsumptionStatus!=null and lastConsumptionStatus==3">
- and (cas.latest_record_state=#{lastConsumptionStatus,
jdbcType=INTEGER}
+ <if test="lastConsumptionStatus != null and lastConsumptionStatus ==
3">
+ and (cas.latest_record_state = #{lastConsumptionStatus,
jdbcType=INTEGER}
or cas.latest_record_state is null)
</if>
order by id desc
@@ -321,15 +292,8 @@
select status as `key`, count(1) as value
from consumption
where is_deleted=0
- <if test="consumerGroupName != null and consumerGroupName != ''">
- and consumer_group_name = #{consumerGroupName,jdbcType=VARCHAR}
- </if>
- <if test="consumerGroupId != null and consumerGroupId != ''">
- and consumer_group_id = #{consumerGroupId,jdbcType=VARCHAR}
- </if>
- <if test="consumerGroupIdLike != null and consumerGroupIdLike != ''">
- <bind name="consumerGroupIdLike" value="'%' +
_parameter.consumerGroupIdLike + '%'"/>
- and consumer_group_id like #{consumerGroupIdLike,jdbcType=VARCHAR}
+ <if test="consumerGroup != null and consumerGroup != ''">
+ and consumer_group = #{consumerGroup,jdbcType=VARCHAR}
</if>
<if test="inCharges != null and inCharges != ''">
and FIND_IN_SET(#{inCharges,jdbcType=VARCHAR},in_charges)
@@ -337,16 +301,12 @@
<if test="inlongGroupId != null and inlongGroupId != ''">
and inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR}
</if>
- <if test="middlewareType != null and middlewareType != ''">
- and middleware_type = #{middlewareType,jdbcType=VARCHAR}
+ <if test="mqType != null and mqType != ''">
+ and mq_type = #{mqType,jdbcType=VARCHAR}
</if>
<if test="topic != null and topic != ''">
and topic = #{topic,jdbcType=VARCHAR}
</if>
- <if test="topicLike != null and topicLike != ''">
- <bind name="topicLike" value="'%' + _parameter.topicLike + '%'"/>
- and topic like #{topicLike,jdbcType=VARCHAR}
- </if>
<if test="filterEnabled != null">
and filter_enabled = #{filterEnabled,jdbcType=INTEGER}
</if>
@@ -365,16 +325,11 @@
<if test="userName != null and userName !=''">
and (
FIND_IN_SET(#{userName,jdbcType=VARCHAR},in_charges)
- or
- creator = #{userName,jdbcType=VARCHAR}
+ or creator = #{userName,jdbcType=VARCHAR}
)
</if>
<if test="keyword != null and keyword !=''">
- <bind name="keyword" value="'%' + _parameter.keyword + '%'"/>
- and (
- topic like #{keyword,jdbcType=VARCHAR}
- or consumer_group_id like #{keyword,jdbcType=VARCHAR}
- )
+ and ( topic like CONCAT('%', #{keyword}, '%') or consumer_group
like CONCAT('%', #{keyword}, '%') )
</if>
group by status
</select>
diff --git
a/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionPulsarEntityMapper.xml
b/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionPulsarEntityMapper.xml
index ed4489403..a74fdb333 100644
---
a/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionPulsarEntityMapper.xml
+++
b/inlong-manager/manager-dao/src/main/resources/mappers/ConsumptionPulsarEntityMapper.xml
@@ -22,7 +22,7 @@
<resultMap id="BaseResultMap"
type="org.apache.inlong.manager.dao.entity.ConsumptionPulsarEntity">
<id column="id" jdbcType="INTEGER" property="id"/>
<result column="consumption_id" jdbcType="VARCHAR"
property="consumptionId"/>
- <result column="consumer_group_id" jdbcType="VARCHAR"
property="consumerGroupId"/>
+ <result column="consumer_group" jdbcType="VARCHAR"
property="consumerGroup"/>
<result column="inlong_group_id" jdbcType="VARCHAR"
property="inlongGroupId"/>
<result column="is_rlq" jdbcType="INTEGER" property="isRlq"/>
<result column="retry_letter_topic" jdbcType="VARCHAR"
property="retryLetterTopic"/>
@@ -31,7 +31,7 @@
<result column="is_deleted" jdbcType="INTEGER" property="isDeleted"/>
</resultMap>
<sql id="Base_Column_List">
- id, consumption_id, consumer_group_id, inlong_group_id, is_rlq,
retry_letter_topic,
+ id, consumption_id, consumer_group, inlong_group_id, is_rlq,
retry_letter_topic,
is_dlq, dead_letter_topic, is_deleted
</sql>
<select id="selectByPrimaryKey" parameterType="java.lang.Integer"
resultMap="BaseResultMap">
@@ -56,10 +56,10 @@
</delete>
<insert id="insert"
parameterType="org.apache.inlong.manager.dao.entity.ConsumptionPulsarEntity">
- insert into consumption_pulsar (id, consumption_id, consumer_group_id,
+ insert into consumption_pulsar (id, consumption_id, consumer_group,
inlong_group_id, is_rlq,
retry_letter_topic,
is_dlq, dead_letter_topic, is_deleted)
- values (#{id,jdbcType=INTEGER}, #{consumptionId,jdbcType=INTEGER},
#{consumerGroupId,jdbcType=VARCHAR},
+ values (#{id,jdbcType=INTEGER}, #{consumptionId,jdbcType=INTEGER},
#{consumerGroup,jdbcType=VARCHAR},
#{inlongGroupId,jdbcType=VARCHAR}, #{isRlq,jdbcType=INTEGER},
#{retryLetterTopic,jdbcType=VARCHAR},
#{isDlq,jdbcType=INTEGER},
#{deadLetterTopic,jdbcType=VARCHAR}, #{isDeleted,jdbcType=INTEGER})
</insert>
@@ -72,8 +72,8 @@
<if test="consumptionId != null">
consumption_id,
</if>
- <if test="consumerGroupId != null">
- consumer_group_id,
+ <if test="consumerGroup != null">
+ consumer_group,
</if>
<if test="inlongGroupId != null">
inlong_group_id,
@@ -98,8 +98,8 @@
<if test="id != null">
#{id,jdbcType=INTEGER},
</if>
- <if test="consumerGroupId != null">
- #{consumerGroupId,jdbcType=VARCHAR},
+ <if test="consumerGroup != null">
+ #{consumerGroup,jdbcType=VARCHAR},
</if>
<if test="inlongGroupId != null">
#{inlongGroupId,jdbcType=VARCHAR},
@@ -128,8 +128,8 @@
<if test="consumptionId != null">
consumption_id = #{consumptionId,jdbcType=INTEGER},
</if>
- <if test="consumerGroupId != null">
- consumer_group_id = #{consumerGroupId,jdbcType=VARCHAR},
+ <if test="consumerGroup != null">
+ consumer_group = #{consumerGroup,jdbcType=VARCHAR},
</if>
<if test="inlongGroupId != null">
inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR},
@@ -155,7 +155,7 @@
<update id="updateByConsumptionId">
update consumption_pulsar
<set>
- consumer_group_id = #{consumerGroupId,jdbcType=VARCHAR},
+ consumer_group = #{consumerGroup,jdbcType=VARCHAR},
inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR},
is_rlq = #{isRlq,jdbcType=INTEGER},
retry_letter_topic = #{retryLetterTopic,jdbcType=VARCHAR},
@@ -168,7 +168,7 @@
<update id="updateByPrimaryKey"
parameterType="org.apache.inlong.manager.dao.entity.ConsumptionPulsarEntity">
update consumption_pulsar
set consumption_id = #{consumptionId,jdbcType=INTEGER},
- consumer_group_id = #{consumerGroupId,jdbcType=VARCHAR},
+ consumer_group = #{consumerGroup,jdbcType=VARCHAR},
inlong_group_id = #{inlongGroupId,jdbcType=VARCHAR},
is_rlq = #{isRlq,jdbcType=INTEGER},
retry_letter_topic = #{retryLetterTopic,jdbcType=VARCHAR},
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
index 058e12b01..bdad67949 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/ConsumptionService.java
@@ -55,13 +55,13 @@ public interface ConsumptionService {
ConsumptionInfo get(Integer id);
/**
- * Determine whether the Consumer group ID already exists
+ * Determine whether the Consumer group already exists
*
- * @param consumerGroupId Consumer group ID
+ * @param consumerGroup Consumer group
* @param excludeSelfId Exclude the ID of this record
* @return does it exist
*/
- boolean isConsumerGroupIdExists(String consumerGroupId, Integer
excludeSelfId);
+ boolean isConsumerGroupExists(String consumerGroup, Integer excludeSelfId);
/**
* Save basic data consumption information
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
index 476d633d6..6141be364 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceImpl.java
@@ -124,7 +124,7 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
ConsumptionInfo info = CommonBeanUtils.copyProperties(entity,
ConsumptionInfo::new);
- MQType mqType = MQType.forType(info.getMiddlewareType());
+ MQType mqType = MQType.forType(info.getMqType());
if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
ConsumptionPulsarEntity pulsarEntity =
consumptionPulsarMapper.selectByConsumptionId(info.getId());
Preconditions.checkNotNull(pulsarEntity, "Pulsar consumption
cannot be empty, as the middleware is Pulsar");
@@ -138,9 +138,9 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
}
@Override
- public boolean isConsumerGroupIdExists(String consumerGroup, Integer
excludeSelfId) {
+ public boolean isConsumerGroupExists(String consumerGroup, Integer
excludeSelfId) {
ConsumptionQuery consumptionQuery = new ConsumptionQuery();
- consumptionQuery.setConsumerGroupId(consumerGroup);
+ consumptionQuery.setConsumerGroup(consumerGroup);
consumptionQuery.setIsAdminRole(true);
List<ConsumptionEntity> result =
consumptionMapper.listByQuery(consumptionQuery);
if (excludeSelfId != null) {
@@ -156,7 +156,7 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
fullConsumptionInfo(info);
Date now = new Date();
ConsumptionEntity entity = this.saveConsumption(info, operator, now);
- MQType mqType = MQType.forType(entity.getMiddlewareType());
+ MQType mqType = MQType.forType(entity.getMqType());
if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
savePulsarInfo(info.getMqExtInfo(), entity);
}
@@ -213,7 +213,7 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
Integer consumptionId = entity.getId();
pulsar.setConsumptionId(consumptionId);
pulsar.setInlongGroupId(groupId);
- pulsar.setConsumerGroupId(entity.getConsumerGroupId());
+ pulsar.setConsumerGroup(entity.getConsumerGroup());
pulsar.setIsDeleted(0);
// Pulsar consumer information may already exist, update if it exists,
add if it does not exist
@@ -245,11 +245,11 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
entity.setModifyTime(now);
// Modify Pulsar consumption info
- MQType mqType = MQType.forType(info.getMiddlewareType());
+ MQType mqType = MQType.forType(info.getMqType());
if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
ConsumptionPulsarEntity pulsarEntity =
consumptionPulsarMapper.selectByConsumptionId(consumptionId);
Preconditions.checkNotNull(pulsarEntity, "Pulsar consumption
cannot be null");
- pulsarEntity.setConsumerGroupId(info.getConsumerGroupId());
+ pulsarEntity.setConsumerGroup(info.getConsumerGroup());
// Whether DLQ / RLQ is turned on or off
ConsumptionPulsarInfo update = (ConsumptionPulsarInfo)
info.getMqExtInfo();
@@ -355,10 +355,9 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
MQType mqType = MQType.forType(groupInfo.getMqType());
ConsumptionEntity entity = new ConsumptionEntity();
entity.setInlongGroupId(groupId);
- entity.setMiddlewareType(mqType.getType());
+ entity.setMqType(mqType.getType());
entity.setTopic(topic);
- entity.setConsumerGroupId(consumerGroup);
- entity.setConsumerGroupName(consumerGroup);
+ entity.setConsumerGroup(consumerGroup);
entity.setInCharges(groupInfo.getInCharges());
entity.setFilterEnabled(0);
@@ -372,7 +371,7 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
ConsumptionPulsarEntity pulsarEntity = new
ConsumptionPulsarEntity();
pulsarEntity.setConsumptionId(entity.getId());
- pulsarEntity.setConsumerGroupId(consumerGroup);
+ pulsarEntity.setConsumerGroup(consumerGroup);
pulsarEntity.setInlongGroupId(groupId);
pulsarEntity.setIsDeleted(GlobalConstants.UN_DELETED);
consumptionPulsarMapper.insert(pulsarEntity);
@@ -384,7 +383,7 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
private NewConsumptionProcessForm
genNewConsumptionProcessForm(ConsumptionInfo consumptionInfo) {
NewConsumptionProcessForm form = new NewConsumptionProcessForm();
Integer id = consumptionInfo.getId();
- MQType mqType = MQType.forType(consumptionInfo.getMiddlewareType());
+ MQType mqType = MQType.forType(consumptionInfo.getMqType());
if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
ConsumptionPulsarEntity consumptionPulsarEntity =
consumptionPulsarMapper.selectByConsumptionId(id);
ConsumptionPulsarInfo pulsarInfo =
CommonBeanUtils.copyProperties(consumptionPulsarEntity,
@@ -419,10 +418,8 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
*/
private void fullConsumptionInfo(ConsumptionInfo info) {
Preconditions.checkNotNull(info, "consumption info cannot be null");
- info.setConsumerGroupId(info.getConsumerGroupName());
-
-
Preconditions.checkFalse(isConsumerGroupIdExists(info.getConsumerGroupId(),
info.getId()),
- "consumer group " + info.getConsumerGroupId() + " already
exist");
+
Preconditions.checkFalse(isConsumerGroupExists(info.getConsumerGroup(),
info.getId()),
+ "consumer group " + info.getConsumerGroup() + " already
exist");
if (info.getId() != null) {
ConsumptionEntity consumptionEntity =
consumptionMapper.selectByPrimaryKey(info.getId());
@@ -456,7 +453,7 @@ public class ConsumptionServiceImpl implements
ConsumptionService {
Preconditions.checkEmpty(topicSet, "topic [" + topicSet + "]
not belong to inlong group " + groupId);
}
}
- info.setMiddlewareType(mqType.getType());
+ info.setMqType(mqType.getType());
}
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/mq/MiddlewareFactory.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/mq/MiddlewareFactory.java
index c67acd0ee..4fb948832 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/mq/MiddlewareFactory.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/core/mq/MiddlewareFactory.java
@@ -46,8 +46,7 @@ public class MiddlewareFactory {
return middleware;
}
}
- throw new BusinessException(ErrorCodeEnum.MQ_TYPE_NOT_SUPPORTED,
- "Current version of InLong not support middleware type of MQ:"
+ type.name());
+ throw new BusinessException(ErrorCodeEnum.MQ_TYPE_NOT_SUPPORTED,
"Unsupported MQ type of " + type.name());
}
}
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionCompleteProcessListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionCompleteProcessListener.java
index 4720404ab..19e2a6a26 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionCompleteProcessListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionCompleteProcessListener.java
@@ -84,28 +84,28 @@ public class ConsumptionCompleteProcessListener implements
ProcessEventListener
throw new WorkflowListenerException("consumption not exits for
id=" + consumptionId);
}
- MQType mqType = MQType.forType(entity.getMiddlewareType());
+ MQType mqType = MQType.forType(entity.getMqType());
if (mqType == MQType.TUBE) {
this.createTubeConsumerGroup(entity);
return ListenerResult.success("Create Tube consumer group
successful");
} else if (mqType == MQType.PULSAR || mqType == MQType.TDMQ_PULSAR) {
this.createPulsarTopicMessage(entity);
} else {
- throw new WorkflowListenerException("middleware type [" + mqType +
"] not supported");
+ throw new WorkflowListenerException("Unsupported MQ type [" +
mqType + "]");
}
- this.updateConsumerInfo(consumptionId, entity.getConsumerGroupId());
- return ListenerResult.success("create Tube /Pulsar consumer group
successful");
+ this.updateConsumerInfo(consumptionId, entity.getConsumerGroup());
+ return ListenerResult.success("Create MQ consumer group successful");
}
/**
* Update consumption after approve
*/
- private void updateConsumerInfo(Integer consumptionId, String
consumerGroupId) {
+ private void updateConsumerInfo(Integer consumptionId, String
consumerGroup) {
ConsumptionEntity update = new ConsumptionEntity();
update.setId(consumptionId);
update.setStatus(ConsumptionStatus.APPROVED.getStatus());
- update.setConsumerGroupId(consumerGroupId);
+ update.setConsumerGroup(consumerGroup);
update.setModifyTime(new Date());
consumptionMapper.updateByPrimaryKeySelective(update);
}
@@ -119,7 +119,7 @@ public class ConsumptionCompleteProcessListener implements
ProcessEventListener
Preconditions.checkNotNull(groupInfo, "inlong group not found for
groupId=" + groupId);
String mqResource = groupInfo.getMqResource();
Preconditions.checkNotNull(mqResource, "mq resource cannot empty for
groupId=" + groupId);
- PulsarClusterInfo globalCluster =
commonOperateService.getPulsarClusterInfo(entity.getMiddlewareType());
+ PulsarClusterInfo globalCluster =
commonOperateService.getPulsarClusterInfo(entity.getMqType());
try (PulsarAdmin pulsarAdmin =
PulsarUtils.getPulsarAdmin(globalCluster)) {
PulsarTopicBean topicMessage = new PulsarTopicBean();
String tenant = clusterBean.getDefaultTenant();
@@ -127,7 +127,7 @@ public class ConsumptionCompleteProcessListener implements
ProcessEventListener
topicMessage.setNamespace(mqResource);
// If cross-regional replication is started, each cluster needs to
create consumer groups in cycles
- String consumerGroup = entity.getConsumerGroupId();
+ String consumerGroup = entity.getConsumerGroup();
List<String> clusters = PulsarUtils.getPulsarClusters(pulsarAdmin);
List<String> topics = Arrays.asList(entity.getTopic().split(","));
this.createPulsarSubscription(pulsarAdmin, consumerGroup,
topicMessage, clusters, topics, globalCluster);
@@ -164,7 +164,7 @@ public class ConsumptionCompleteProcessListener implements
ProcessEventListener
addTubeConsumeGroupRequest.setCreateUser(consumption.getCreator());
AddTubeConsumeGroupRequest.GroupNameJsonSetBean bean = new
AddTubeConsumeGroupRequest.GroupNameJsonSetBean();
bean.setTopicName(consumption.getTopic());
- bean.setGroupName(consumption.getConsumerGroupId());
+ bean.setGroupName(consumption.getConsumerGroup());
addTubeConsumeGroupRequest.setGroupNameJsonSet(Collections.singletonList(bean));
try {
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionPassTaskListener.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionPassTaskListener.java
index f9edd90d2..58ee04afe 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionPassTaskListener.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/workflow/consumption/listener/ConsumptionPassTaskListener.java
@@ -23,9 +23,9 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.common.exceptions.WorkflowListenerException;
import org.apache.inlong.manager.common.pojo.consumption.ConsumptionInfo;
-import org.apache.inlong.manager.service.core.ConsumptionService;
import
org.apache.inlong.manager.common.pojo.workflow.form.ConsumptionApproveForm;
import
org.apache.inlong.manager.common.pojo.workflow.form.NewConsumptionProcessForm;
+import org.apache.inlong.manager.service.core.ConsumptionService;
import org.apache.inlong.manager.workflow.WorkflowContext;
import org.apache.inlong.manager.workflow.event.ListenerResult;
import org.apache.inlong.manager.workflow.event.task.TaskEvent;
@@ -53,16 +53,16 @@ public class ConsumptionPassTaskListener implements
TaskEventListener {
NewConsumptionProcessForm form = (NewConsumptionProcessForm)
context.getProcessForm();
ConsumptionApproveForm approveForm = (ConsumptionApproveForm)
context.getActionContext().getForm();
ConsumptionInfo info = form.getConsumptionInfo();
- if (StringUtils.equals(approveForm.getConsumerGroupId(),
info.getConsumerGroupId())) {
- return ListenerResult.success("The consumer group name has not
been modified");
+ if (StringUtils.equals(approveForm.getConsumerGroup(),
info.getConsumerGroup())) {
+ return ListenerResult.success("The consumer group has not been
modified");
}
- boolean exist =
consumptionService.isConsumerGroupIdExists(approveForm.getConsumerGroupId(),
info.getId());
+ boolean exist =
consumptionService.isConsumerGroupExists(approveForm.getConsumerGroup(),
info.getId());
if (exist) {
- log.error("consumerGroupId already exist! duplicate :{}",
approveForm.getConsumerGroupId());
- throw new
BusinessException(ErrorCodeEnum.CONSUMER_GROUP_NAME_DUPLICATED);
+ log.error("consumer group {} already exist",
approveForm.getConsumerGroup());
+ throw new
BusinessException(ErrorCodeEnum.CONSUMER_GROUP_DUPLICATED);
}
- return ListenerResult.success("Consumer group name from" +
info.getConsumerGroupId()
- + "change to " + approveForm.getConsumerGroupId());
+ return ListenerResult.success("Consumer group from " +
info.getConsumerGroup()
+ + " change to " + approveForm.getConsumerGroup());
}
@Override
diff --git
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceTest.java
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceTest.java
index dcb9bd346..6a1ec0698 100644
---
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceTest.java
+++
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/ConsumptionServiceTest.java
@@ -42,14 +42,14 @@ public class ConsumptionServiceTest extends ServiceBaseTest
{
private Integer saveConsumption(String inlongGroup, String consumerGroup,
String operator) {
ConsumptionInfo consumptionInfo = new ConsumptionInfo();
consumptionInfo.setTopic(inlongGroup);
- consumptionInfo.setConsumerGroupName(consumerGroup);
+ consumptionInfo.setConsumerGroup(consumerGroup);
consumptionInfo.setInlongGroupId("b_" + inlongGroup);
- consumptionInfo.setMiddlewareType(MQType.PULSAR.getType());
+ consumptionInfo.setMqType(MQType.PULSAR.getType());
consumptionInfo.setCreator(operator);
consumptionInfo.setInCharges("admin");
ConsumptionPulsarInfo pulsarInfo = new ConsumptionPulsarInfo();
- pulsarInfo.setMiddlewareType(MQType.PULSAR.getType());
+ pulsarInfo.setMqType(MQType.PULSAR.getType());
pulsarInfo.setIsDlq(1);
pulsarInfo.setDeadLetterTopic("test_dlq");
pulsarInfo.setIsRlq(0);
diff --git
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
index 7497d2e0f..7a26ec2cd 100644
---
a/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
+++
b/inlong-manager/manager-service/src/test/java/org/apache/inlong/manager/service/core/impl/InlongGroupProcessOperationTest.java
@@ -56,7 +56,7 @@ public class InlongGroupProcessOperationTest extends
ServiceBaseTest {
private ServiceTaskListenerFactory serviceTaskListenerFactory;
/**
- * Set some base infomation before satart process.
+ * Set some base information before start process.
*/
public void before() {
MockPlugin mockPlugin = new MockPlugin();
@@ -119,7 +119,6 @@ public class InlongGroupProcessOperationTest extends
ServiceBaseTest {
// testRestartProcess();
// boolean result = groupProcessOperation.deleteProcess(GROUP_ID,
OPERATOR);
// Assert.assertTrue(result);
-
}
}
diff --git
a/inlong-manager/manager-test/src/main/resources/sql/apache_inlong_manager.sql
b/inlong-manager/manager-test/src/main/resources/sql/apache_inlong_manager.sql
index ed3847efb..d9a0236c7 100644
---
a/inlong-manager/manager-test/src/main/resources/sql/apache_inlong_manager.sql
+++
b/inlong-manager/manager-test/src/main/resources/sql/apache_inlong_manager.sql
@@ -249,21 +249,20 @@ CREATE TABLE IF NOT EXISTS `common_file_server`
-- ----------------------------
CREATE TABLE IF NOT EXISTS `consumption`
(
- `id` int(11) NOT NULL AUTO_INCREMENT COMMENT
'Incremental primary key',
- `consumer_group_name` varchar(256) DEFAULT NULL COMMENT 'consumer
group name',
- `consumer_group_id` varchar(256) NOT NULL COMMENT 'Consumer group ID',
- `in_charges` varchar(512) NOT NULL COMMENT 'Person in charge of
consumption',
- `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group id',
- `middleware_type` varchar(10) DEFAULT 'TUBE' COMMENT 'The
middleware type of message queue, high throughput: TUBE, high consistency:
PULSAR',
- `topic` varchar(256) NOT NULL COMMENT 'Consumption topic',
- `filter_enabled` int(2) DEFAULT '0' COMMENT 'Whether to
filter, default 0, not filter consume',
- `inlong_stream_id` varchar(256) DEFAULT NULL COMMENT 'Inlong
stream ID for consumption, if filter_enable is 1, it cannot empty',
- `status` int(4) NOT NULL COMMENT 'Status: draft,
pending approval, approval rejected, approval passed',
- `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to
delete, 0: not deleted, > 0: deleted',
- `creator` varchar(64) NOT NULL COMMENT 'creator',
- `modifier` varchar(64) DEFAULT NULL COMMENT 'modifier',
- `create_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP COMMENT
'Create time',
- `modify_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP ON
UPDATE CURRENT_TIMESTAMP COMMENT 'Modify time',
+ `id` int(11) NOT NULL AUTO_INCREMENT COMMENT
'Incremental primary key',
+ `consumer_group` varchar(256) NOT NULL COMMENT 'Consumer group',
+ `in_charges` varchar(512) NOT NULL COMMENT 'Person in charge of
consumption',
+ `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group id',
+ `mq_type` varchar(10) DEFAULT 'TUBE' COMMENT 'Message queue
type, high throughput: TUBE, high consistency: PULSAR',
+ `topic` varchar(256) NOT NULL COMMENT 'Consumption topic',
+ `filter_enabled` int(2) DEFAULT '0' COMMENT 'Whether to
filter, default 0, not filter consume',
+ `inlong_stream_id` varchar(256) DEFAULT NULL COMMENT 'Inlong stream
ID for consumption, if filter_enable is 1, it cannot empty',
+ `status` int(4) NOT NULL COMMENT 'Status: draft, pending
approval, approval rejected, approval passed',
+ `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to
delete, 0: not deleted, > 0: deleted',
+ `creator` varchar(64) NOT NULL COMMENT 'creator',
+ `modifier` varchar(64) DEFAULT NULL COMMENT 'modifier',
+ `create_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP COMMENT
'Create time',
+ `modify_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE
CURRENT_TIMESTAMP COMMENT 'Modify time',
PRIMARY KEY (`id`)
);
@@ -272,16 +271,15 @@ CREATE TABLE IF NOT EXISTS `consumption`
-- ----------------------------
CREATE TABLE IF NOT EXISTS `consumption_pulsar`
(
- `id` int(11) NOT NULL AUTO_INCREMENT,
- `consumption_id` int(11) DEFAULT NULL COMMENT 'ID of the
consumption information to which it belongs, guaranteed to be uniquely
associated with consumption information',
- `consumer_group_id` varchar(256) NOT NULL COMMENT 'Consumer group ID',
- `consumer_group_name` varchar(256) DEFAULT NULL COMMENT 'Consumer group
name',
- `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group ID',
- `is_rlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure the retry letter topic, 0: no configuration, 1: configuration',
- `retry_letter_topic` varchar(256) DEFAULT NULL COMMENT 'The name of the
retry queue topic',
- `is_dlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure dead letter topic, 0: no configuration, 1: means configuration',
- `dead_letter_topic` varchar(256) DEFAULT NULL COMMENT 'dead letter topic
name',
- `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to delete,
0: not deleted, > 0: deleted',
+ `id` int(11) NOT NULL AUTO_INCREMENT,
+ `consumption_id` int(11) DEFAULT NULL COMMENT 'ID of the
consumption information to which it belongs, guaranteed to be uniquely
associated with consumption information',
+ `consumer_group` varchar(256) NOT NULL COMMENT 'Consumer group',
+ `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group ID',
+ `is_rlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure the retry letter topic, 0: no configuration, 1: configuration',
+ `retry_letter_topic` varchar(256) DEFAULT NULL COMMENT 'The name of the
retry queue topic',
+ `is_dlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure dead letter topic, 0: no configuration, 1: means configuration',
+ `dead_letter_topic` varchar(256) DEFAULT NULL COMMENT 'dead letter topic
name',
+ `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to delete,
0: not deleted, > 0: deleted',
PRIMARY KEY (`id`)
) COMMENT ='Pulsar consumption table';
diff --git a/inlong-manager/manager-web/sql/apache_inlong_manager.sql
b/inlong-manager/manager-web/sql/apache_inlong_manager.sql
index dc1ccaa47..eab1c7927 100644
--- a/inlong-manager/manager-web/sql/apache_inlong_manager.sql
+++ b/inlong-manager/manager-web/sql/apache_inlong_manager.sql
@@ -264,21 +264,20 @@ CREATE TABLE IF NOT EXISTS `common_file_server`
-- ----------------------------
CREATE TABLE IF NOT EXISTS `consumption`
(
- `id` int(11) NOT NULL AUTO_INCREMENT COMMENT
'Incremental primary key',
- `consumer_group_name` varchar(256) DEFAULT NULL COMMENT 'consumer
group name',
- `consumer_group_id` varchar(256) NOT NULL COMMENT 'Consumer group ID',
- `in_charges` varchar(512) NOT NULL COMMENT 'Person in charge of
consumption',
- `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group id',
- `middleware_type` varchar(10) DEFAULT 'TUBE' COMMENT 'The
middleware type of message queue, high throughput: TUBE, high consistency:
PULSAR',
- `topic` varchar(256) NOT NULL COMMENT 'Consumption topic',
- `filter_enabled` int(2) DEFAULT '0' COMMENT 'Whether to
filter, default 0, not filter consume',
- `inlong_stream_id` varchar(256) DEFAULT NULL COMMENT 'Inlong
stream ID for consumption, if filter_enable is 1, it cannot empty',
- `status` int(4) NOT NULL COMMENT 'Status: draft,
pending approval, approval rejected, approval passed',
- `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to
delete, 0: not deleted, > 0: deleted',
- `creator` varchar(64) NOT NULL COMMENT 'creator',
- `modifier` varchar(64) DEFAULT NULL COMMENT 'modifier',
- `create_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP COMMENT
'Create time',
- `modify_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP ON
UPDATE CURRENT_TIMESTAMP COMMENT 'Modify time',
+ `id` int(11) NOT NULL AUTO_INCREMENT COMMENT
'Incremental primary key',
+ `consumer_group` varchar(256) NOT NULL COMMENT 'Consumer group',
+ `in_charges` varchar(512) NOT NULL COMMENT 'Person in charge of
consumption',
+ `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group id',
+ `mq_type` varchar(10) DEFAULT 'TUBE' COMMENT 'Message
queue, high throughput: TUBE, high consistency: PULSAR',
+ `topic` varchar(256) NOT NULL COMMENT 'Consumption topic',
+ `filter_enabled` int(2) DEFAULT '0' COMMENT 'Whether to
filter, default 0, not filter consume',
+ `inlong_stream_id` varchar(256) DEFAULT NULL COMMENT 'Inlong stream
ID for consumption, if filter_enable is 1, it cannot empty',
+ `status` int(4) NOT NULL COMMENT 'Status: draft, pending
approval, approval rejected, approval passed',
+ `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to
delete, 0: not deleted, > 0: deleted',
+ `creator` varchar(64) NOT NULL COMMENT 'creator',
+ `modifier` varchar(64) DEFAULT NULL COMMENT 'modifier',
+ `create_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP COMMENT
'Create time',
+ `modify_time` timestamp NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE
CURRENT_TIMESTAMP COMMENT 'Modify time',
PRIMARY KEY (`id`)
) ENGINE = InnoDB
DEFAULT CHARSET = utf8mb4 COMMENT ='Data consumption configuration table';
@@ -288,16 +287,15 @@ CREATE TABLE IF NOT EXISTS `consumption`
-- ----------------------------
CREATE TABLE IF NOT EXISTS `consumption_pulsar`
(
- `id` int(11) NOT NULL AUTO_INCREMENT,
- `consumption_id` int(11) DEFAULT NULL COMMENT 'ID of the
consumption information to which it belongs, guaranteed to be uniquely
associated with consumption information',
- `consumer_group_id` varchar(256) NOT NULL COMMENT 'Consumer group ID',
- `consumer_group_name` varchar(256) DEFAULT NULL COMMENT 'Consumer group
name',
- `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group ID',
- `is_rlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure the retry letter topic, 0: no configuration, 1: configuration',
- `retry_letter_topic` varchar(256) DEFAULT NULL COMMENT 'The name of the
retry queue topic',
- `is_dlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure dead letter topic, 0: no configuration, 1: means configuration',
- `dead_letter_topic` varchar(256) DEFAULT NULL COMMENT 'dead letter topic
name',
- `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to delete,
0: not deleted, > 0: deleted',
+ `id` int(11) NOT NULL AUTO_INCREMENT,
+ `consumption_id` int(11) DEFAULT NULL COMMENT 'ID of the
consumption information to which it belongs, guaranteed to be uniquely
associated with consumption information',
+ `consumer_group` varchar(256) NOT NULL COMMENT 'Consumer group',
+ `inlong_group_id` varchar(256) NOT NULL COMMENT 'Inlong group ID',
+ `is_rlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure the retry letter topic, 0: no configuration, 1: configuration',
+ `retry_letter_topic` varchar(256) DEFAULT NULL COMMENT 'The name of the
retry queue topic',
+ `is_dlq` tinyint(1) DEFAULT '0' COMMENT 'Whether to
configure dead letter topic, 0: no configuration, 1: means configuration',
+ `dead_letter_topic` varchar(256) DEFAULT NULL COMMENT 'dead letter topic
name',
+ `is_deleted` int(11) DEFAULT '0' COMMENT 'Whether to delete,
0: not deleted, > 0: deleted',
PRIMARY KEY (`id`)
) ENGINE = InnoDB
DEFAULT CHARSET = utf8mb4 COMMENT ='Pulsar consumption table';