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 3675d3d73 [INLONG-7701][Manager] Support test connection for hudi in 
NodeManagement (#7710)
3675d3d73 is described below

commit 3675d3d734c0d349d65cd934ac842fbc2f7f19c0
Author: feat <[email protected]>
AuthorDate: Tue Mar 28 15:56:20 2023 +0800

    [INLONG-7701][Manager] Support test connection for hudi in NodeManagement 
(#7710)
---
 .../service/node/hudi/HudiDataNodeOperator.java    | 22 ++++++++++++++++++++++
 .../node/iceberg/IcebergDataNodeOperator.java      |  2 +-
 .../resource/sink/hudi/HudiCatalogClient.java      | 15 +++++++++++++--
 3 files changed, 36 insertions(+), 3 deletions(-)

diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
index 3d8e68cfe..e677b5090 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
@@ -23,6 +23,7 @@ import org.apache.inlong.manager.common.consts.DataNodeType;
 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.Preconditions;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeRequest;
@@ -30,6 +31,7 @@ import 
org.apache.inlong.manager.pojo.node.hudi.HudiDataNodeDTO;
 import org.apache.inlong.manager.pojo.node.hudi.HudiDataNodeInfo;
 import org.apache.inlong.manager.pojo.node.hudi.HudiDataNodeRequest;
 import org.apache.inlong.manager.service.node.AbstractDataNodeOperator;
+import org.apache.inlong.manager.service.resource.sink.hudi.HudiCatalogClient;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 import org.springframework.beans.factory.annotation.Autowired;
@@ -81,4 +83,24 @@ public class HudiDataNodeOperator extends 
AbstractDataNodeOperator {
         }
     }
 
+    @Override
+    public Boolean testConnection(DataNodeRequest request) {
+        HudiDataNodeRequest hudiRequest = (HudiDataNodeRequest) request;
+        String metastoreUri = hudiRequest.getUrl();
+        String warehouse = hudiRequest.getWarehouse();
+        Preconditions.expectNotBlank(metastoreUri, 
ErrorCodeEnum.INVALID_PARAMETER, "connection url cannot be empty");
+        try (HudiCatalogClient client = new HudiCatalogClient(metastoreUri, 
warehouse)) {
+            client.open();
+            client.listAllDatabases();
+            LOGGER.info("hudi connection not null - connection success for 
metastoreUri={}, warehouse={}",
+                    metastoreUri, warehouse);
+            return true;
+        } catch (Exception e) {
+            String errMsg = String.format("hudi connection failed for 
metastoreUri=%s, warehouse=%s", metastoreUri,
+                    warehouse);
+            LOGGER.error(errMsg, e);
+            throw new BusinessException(errMsg);
+        }
+    }
+
 }
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
index ece8bf300..7c585df52 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
@@ -96,7 +96,7 @@ public class IcebergDataNodeOperator extends 
AbstractDataNodeOperator {
                     metastoreUri, warehouse);
             return true;
         } catch (Exception e) {
-            String errMsg = String.format("iceberg connection failed for 
metastoreUri=%s, warhouse=%s", metastoreUri,
+            String errMsg = String.format("iceberg connection failed for 
metastoreUri=%s, warehouse=%s", metastoreUri,
                     warehouse);
             LOGGER.error(errMsg, e);
             throw new BusinessException(errMsg);
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/hudi/HudiCatalogClient.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/hudi/HudiCatalogClient.java
index 72eee8360..ede46f8ba 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/hudi/HudiCatalogClient.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/hudi/HudiCatalogClient.java
@@ -50,11 +50,11 @@ import static 
org.apache.inlong.manager.service.resource.sink.hudi.HudiUtils.get
 /**
  * The Catalog client for Hudi.
  */
-public class HudiCatalogClient {
+public class HudiCatalogClient implements AutoCloseable {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(HudiCatalogClient.class);
 
-    private final String dbName;
+    private String dbName;
     private final String warehouse;
     private IMetaStoreClient client;
     private final HiveConf hiveConf;
@@ -67,6 +67,13 @@ public class HudiCatalogClient {
         hiveConf.setBoolVar(HiveConf.ConfVars.METASTORE_EXECUTE_SET_UGI, 
false);
     }
 
+    public HudiCatalogClient(String uri, String warehouse) throws 
MetaException {
+        this.warehouse = warehouse;
+        hiveConf = new HiveConf();
+        hiveConf.setVar(HiveConf.ConfVars.METASTOREURIS, uri);
+        hiveConf.setBoolVar(HiveConf.ConfVars.METASTORE_EXECUTE_SET_UGI, 
false);
+    }
+
     /**
      * Open the hive metastore connection
      */
@@ -227,6 +234,10 @@ public class HudiCatalogClient {
         client.createTable(hiveTable);
     }
 
+    public List<String> listAllDatabases() throws TException {
+        return client.getAllDatabases();
+    }
+
     /**
      * Close the connection of hive metastore
      */

Reply via email to