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 0fa2f05163 [INLONG-9533][Manager] Support setting dataNode when
configuring streamSource for Pulsar, Iceberg and PostgreSQL (#9534)
0fa2f05163 is described below
commit 0fa2f0516360a274e87ed920dd389f00ce6d7b53
Author: fuweng11 <[email protected]>
AuthorDate: Wed Dec 27 21:43:02 2023 +0800
[INLONG-9533][Manager] Support setting dataNode when configuring
streamSource for Pulsar, Iceberg and PostgreSQL (#9534)
---
.../service/source/AbstractSourceOperator.java | 3 ++
.../source/binlog/BinlogSourceOperator.java | 7 ++---
.../source/iceberg/IcebergSourceOperator.java | 29 +++++++++++++++++++
.../postgresql/PostgreSQLSourceOperator.java | 33 ++++++++++++++++++++++
.../source/pulsar/PulsarSourceOperator.java | 25 ++++++++++++++++
5 files changed, 92 insertions(+), 5 deletions(-)
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/AbstractSourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/AbstractSourceOperator.java
index 935a688d82..00f85052fd 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/AbstractSourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/AbstractSourceOperator.java
@@ -35,6 +35,7 @@ import org.apache.inlong.manager.pojo.common.PageResult;
import org.apache.inlong.manager.pojo.source.SourceRequest;
import org.apache.inlong.manager.pojo.source.StreamSource;
import org.apache.inlong.manager.pojo.stream.StreamField;
+import org.apache.inlong.manager.service.node.DataNodeService;
import com.github.pagehelper.Page;
import org.apache.commons.collections.CollectionUtils;
@@ -62,6 +63,8 @@ public abstract class AbstractSourceOperator implements
StreamSourceOperator {
protected StreamSourceFieldEntityMapper sourceFieldMapper;
@Autowired
protected InlongStreamFieldEntityMapper streamFieldMapper;
+ @Autowired
+ protected DataNodeService dataNodeService;
/**
* Getting the source type.
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/binlog/BinlogSourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/binlog/BinlogSourceOperator.java
index 5342d2e56a..32b2256318 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/binlog/BinlogSourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/binlog/BinlogSourceOperator.java
@@ -32,7 +32,6 @@ import
org.apache.inlong.manager.pojo.source.mysql.MySQLBinlogSource;
import org.apache.inlong.manager.pojo.source.mysql.MySQLBinlogSourceDTO;
import org.apache.inlong.manager.pojo.source.mysql.MySQLBinlogSourceRequest;
import org.apache.inlong.manager.pojo.stream.StreamField;
-import org.apache.inlong.manager.service.node.DataNodeOperateHelper;
import org.apache.inlong.manager.service.source.AbstractSourceOperator;
import com.fasterxml.jackson.databind.ObjectMapper;
@@ -49,8 +48,6 @@ import java.util.Objects;
@Service
public class BinlogSourceOperator extends AbstractSourceOperator {
- @Autowired
- protected DataNodeOperateHelper dataNodeHelper;
@Autowired
private ObjectMapper objectMapper;
@@ -69,7 +66,7 @@ public class BinlogSourceOperator extends
AbstractSourceOperator {
MySQLBinlogSourceDTO mySQLBinlogSourceDTO =
JsonUtils.parseObject(sourceEntity.getExtParams(),
MySQLBinlogSourceDTO.class);
if (Objects.nonNull(mySQLBinlogSourceDTO) &&
StringUtils.isBlank(mySQLBinlogSourceDTO.getHostname())) {
- MySQLDataNodeInfo dataNodeInfo = (MySQLDataNodeInfo)
dataNodeHelper.getDataNodeInfo(
+ MySQLDataNodeInfo dataNodeInfo = (MySQLDataNodeInfo)
dataNodeService.get(
sourceEntity.getDataNodeName(), DataNodeType.MYSQL);
CommonBeanUtils.copyProperties(dataNodeInfo, mySQLBinlogSourceDTO,
true);
mySQLBinlogSourceDTO.setUser(dataNodeInfo.getUsername());
@@ -107,7 +104,7 @@ public class BinlogSourceOperator extends
AbstractSourceOperator {
throw new
BusinessException(ErrorCodeEnum.SOURCE_INFO_INCORRECT,
"mysql url and data node is blank");
}
- MySQLDataNodeInfo dataNodeInfo = (MySQLDataNodeInfo)
dataNodeHelper.getDataNodeInfo(
+ MySQLDataNodeInfo dataNodeInfo = (MySQLDataNodeInfo)
dataNodeService.get(
entity.getDataNodeName(), DataNodeType.MYSQL);
CommonBeanUtils.copyProperties(dataNodeInfo, dto, true);
dto.setUser(dataNodeInfo.getUsername());
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/iceberg/IcebergSourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/iceberg/IcebergSourceOperator.java
index d594ab26c0..b333711084 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/iceberg/IcebergSourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/iceberg/IcebergSourceOperator.java
@@ -17,13 +17,16 @@
package org.apache.inlong.manager.service.source.iceberg;
+import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.consts.InlongConstants;
import org.apache.inlong.manager.common.consts.SourceType;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.common.util.CommonBeanUtils;
+import org.apache.inlong.manager.common.util.JsonUtils;
import org.apache.inlong.manager.dao.entity.InlongStreamFieldEntity;
import org.apache.inlong.manager.dao.entity.StreamSourceEntity;
+import org.apache.inlong.manager.pojo.node.iceberg.IcebergDataNodeInfo;
import org.apache.inlong.manager.pojo.sink.iceberg.IcebergColumnInfo;
import org.apache.inlong.manager.pojo.sort.util.FieldInfoUtils;
import org.apache.inlong.manager.pojo.source.SourceRequest;
@@ -37,6 +40,7 @@ import
org.apache.inlong.manager.service.source.AbstractSourceOperator;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.commons.collections.CollectionUtils;
+import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.beans.factory.annotation.Autowired;
@@ -46,6 +50,7 @@ import
org.springframework.transaction.annotation.Transactional;
import java.util.ArrayList;
import java.util.List;
+import java.util.Objects;
/**
* Iceberg stream source operator
@@ -68,6 +73,20 @@ public class IcebergSourceOperator extends
AbstractSourceOperator {
return SourceType.ICEBERG;
}
+ @Override
+ public String getExtParams(StreamSourceEntity sourceEntity) {
+ IcebergSourceDTO icebergSourceDTO =
JsonUtils.parseObject(sourceEntity.getExtParams(),
+ IcebergSourceDTO.class);
+ if (Objects.nonNull(icebergSourceDTO) &&
StringUtils.isBlank(icebergSourceDTO.getUri())) {
+ IcebergDataNodeInfo dataNodeInfo = (IcebergDataNodeInfo)
dataNodeService.get(
+ sourceEntity.getDataNodeName(), DataNodeType.ICEBERG);
+ CommonBeanUtils.copyProperties(dataNodeInfo, icebergSourceDTO,
true);
+ icebergSourceDTO.setUri(dataNodeInfo.getUrl());
+ return JsonUtils.toJsonString(icebergSourceDTO);
+ }
+ return sourceEntity.getExtParams();
+ }
+
@Override
protected void setTargetEntity(SourceRequest request, StreamSourceEntity
targetEntity) {
IcebergSourceRequest sourceRequest = (IcebergSourceRequest) request;
@@ -89,6 +108,16 @@ public class IcebergSourceOperator extends
AbstractSourceOperator {
}
IcebergSourceDTO dto =
IcebergSourceDTO.getFromJson(entity.getExtParams());
+ if (StringUtils.isBlank(dto.getUri())) {
+ if (StringUtils.isBlank(entity.getDataNodeName())) {
+ throw new BusinessException(ErrorCodeEnum.SINK_INFO_INCORRECT,
+ "iceberg catalog uri unspecified and data node is
blank");
+ }
+ IcebergDataNodeInfo dataNodeInfo = (IcebergDataNodeInfo)
dataNodeService.get(
+ entity.getDataNodeName(), DataNodeType.ICEBERG);
+ CommonBeanUtils.copyProperties(dataNodeInfo, dto, true);
+ dto.setUri(dataNodeInfo.getUrl());
+ }
CommonBeanUtils.copyProperties(entity, source, true);
CommonBeanUtils.copyProperties(dto, source, true);
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/postgresql/PostgreSQLSourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/postgresql/PostgreSQLSourceOperator.java
index 161fd4350a..6b360effc7 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/postgresql/PostgreSQLSourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/postgresql/PostgreSQLSourceOperator.java
@@ -17,11 +17,15 @@
package org.apache.inlong.manager.service.source.postgresql;
+import org.apache.inlong.manager.common.consts.DataNodeType;
+import org.apache.inlong.manager.common.consts.InlongConstants;
import org.apache.inlong.manager.common.consts.SourceType;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.common.util.CommonBeanUtils;
+import org.apache.inlong.manager.common.util.JsonUtils;
import org.apache.inlong.manager.dao.entity.StreamSourceEntity;
+import org.apache.inlong.manager.pojo.node.postgresql.PostgreSQLDataNodeInfo;
import org.apache.inlong.manager.pojo.source.SourceRequest;
import org.apache.inlong.manager.pojo.source.StreamSource;
import org.apache.inlong.manager.pojo.source.postgresql.PostgreSQLSource;
@@ -31,6 +35,7 @@ import org.apache.inlong.manager.pojo.stream.StreamField;
import org.apache.inlong.manager.service.source.AbstractSourceOperator;
import com.fasterxml.jackson.databind.ObjectMapper;
+import org.apache.commons.lang3.StringUtils;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Service;
@@ -55,6 +60,22 @@ public class PostgreSQLSourceOperator extends
AbstractSourceOperator {
return SourceType.POSTGRESQL;
}
+ @Override
+ public String getExtParams(StreamSourceEntity sourceEntity) {
+ PostgreSQLSourceDTO postgreSQLSourceDTO =
JsonUtils.parseObject(sourceEntity.getExtParams(),
+ PostgreSQLSourceDTO.class);
+ if (java.util.Objects.nonNull(postgreSQLSourceDTO) &&
StringUtils.isBlank(postgreSQLSourceDTO.getHostname())) {
+ PostgreSQLDataNodeInfo dataNodeInfo = (PostgreSQLDataNodeInfo)
dataNodeService.get(
+ sourceEntity.getDataNodeName(), DataNodeType.POSTGRESQL);
+ CommonBeanUtils.copyProperties(dataNodeInfo, postgreSQLSourceDTO,
true);
+
postgreSQLSourceDTO.setHostname(dataNodeInfo.getUrl().split(InlongConstants.COLON)[0]);
+
postgreSQLSourceDTO.setPort(Integer.valueOf(dataNodeInfo.getUrl().split(InlongConstants.COLON)[1]));
+ postgreSQLSourceDTO.setPassword(dataNodeInfo.getToken());
+ return JsonUtils.toJsonString(postgreSQLSourceDTO);
+ }
+ return sourceEntity.getExtParams();
+ }
+
@Override
protected void setTargetEntity(SourceRequest request, StreamSourceEntity
targetEntity) {
PostgreSQLSourceRequest sourceRequest = (PostgreSQLSourceRequest)
request;
@@ -76,6 +97,18 @@ public class PostgreSQLSourceOperator extends
AbstractSourceOperator {
}
PostgreSQLSourceDTO dto =
PostgreSQLSourceDTO.getFromJson(entity.getExtParams());
+ if (StringUtils.isBlank(dto.getHostname())) {
+ if (StringUtils.isBlank(entity.getDataNodeName())) {
+ throw new BusinessException(ErrorCodeEnum.SINK_INFO_INCORRECT,
+ "postgreSQl hostname unspecified and data node is
blank");
+ }
+ PostgreSQLDataNodeInfo dataNodeInfo = (PostgreSQLDataNodeInfo)
dataNodeService.get(
+ entity.getDataNodeName(), DataNodeType.POSTGRESQL);
+ CommonBeanUtils.copyProperties(dataNodeInfo, dto, true);
+
dto.setHostname(dataNodeInfo.getUrl().split(InlongConstants.COLON)[0]);
+
dto.setPort(Integer.valueOf(dataNodeInfo.getUrl().split(InlongConstants.COLON)[1]));
+ dto.setPassword(dataNodeInfo.getToken());
+ }
CommonBeanUtils.copyProperties(entity, source, true);
CommonBeanUtils.copyProperties(dto, source, true);
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/pulsar/PulsarSourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/pulsar/PulsarSourceOperator.java
index 989031184f..b036db1cb4 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/pulsar/PulsarSourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/source/pulsar/PulsarSourceOperator.java
@@ -18,11 +18,13 @@
package org.apache.inlong.manager.service.source.pulsar;
import org.apache.inlong.common.enums.DataTypeEnum;
+import org.apache.inlong.manager.common.consts.DataNodeType;
import org.apache.inlong.manager.common.consts.SourceType;
import org.apache.inlong.manager.common.enums.ClusterType;
import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
import org.apache.inlong.manager.common.exceptions.BusinessException;
import org.apache.inlong.manager.common.util.CommonBeanUtils;
+import org.apache.inlong.manager.common.util.JsonUtils;
import org.apache.inlong.manager.dao.entity.StreamSinkEntity;
import org.apache.inlong.manager.dao.entity.StreamSourceEntity;
import org.apache.inlong.manager.dao.mapper.StreamSinkEntityMapper;
@@ -30,6 +32,7 @@ import org.apache.inlong.manager.pojo.cluster.ClusterInfo;
import org.apache.inlong.manager.pojo.cluster.pulsar.PulsarClusterInfo;
import org.apache.inlong.manager.pojo.group.InlongGroupInfo;
import org.apache.inlong.manager.pojo.group.pulsar.InlongPulsarInfo;
+import org.apache.inlong.manager.pojo.node.pulsar.PulsarDataNodeInfo;
import org.apache.inlong.manager.pojo.source.SourceRequest;
import org.apache.inlong.manager.pojo.source.StreamSource;
import org.apache.inlong.manager.pojo.source.kafka.KafkaSource;
@@ -84,6 +87,19 @@ public class PulsarSourceOperator extends
AbstractSourceOperator {
return SourceType.PULSAR;
}
+ @Override
+ public String getExtParams(StreamSourceEntity sourceEntity) {
+ PulsarSourceDTO pulsarSourceDTO =
JsonUtils.parseObject(sourceEntity.getExtParams(),
+ PulsarSourceDTO.class);
+ if (java.util.Objects.nonNull(pulsarSourceDTO) &&
StringUtils.isBlank(pulsarSourceDTO.getAdminUrl())) {
+ PulsarDataNodeInfo dataNodeInfo = (PulsarDataNodeInfo)
dataNodeService.get(
+ sourceEntity.getDataNodeName(), DataNodeType.PULSAR);
+ CommonBeanUtils.copyProperties(dataNodeInfo, pulsarSourceDTO,
true);
+ return JsonUtils.toJsonString(pulsarSourceDTO);
+ }
+ return sourceEntity.getExtParams();
+ }
+
@Override
protected void setTargetEntity(SourceRequest request, StreamSourceEntity
targetEntity) {
PulsarSourceRequest sourceRequest = (PulsarSourceRequest) request;
@@ -105,6 +121,15 @@ public class PulsarSourceOperator extends
AbstractSourceOperator {
}
PulsarSourceDTO dto =
PulsarSourceDTO.getFromJson(entity.getExtParams());
+ if (StringUtils.isBlank(dto.getAdminUrl())) {
+ if (StringUtils.isBlank(entity.getDataNodeName())) {
+ throw new BusinessException(ErrorCodeEnum.SINK_INFO_INCORRECT,
+ "pulsar admin url unspecified and data node is blank");
+ }
+ PulsarDataNodeInfo dataNodeInfo = (PulsarDataNodeInfo)
dataNodeService.get(
+ entity.getDataNodeName(), DataNodeType.PULSAR);
+ CommonBeanUtils.copyProperties(dataNodeInfo, dto, true);
+ }
CommonBeanUtils.copyProperties(entity, source, true);
CommonBeanUtils.copyProperties(dto, source, true);