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

Reply via email to