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
*/