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 29242e316 [INLONG-4263][Manager] Support HBase sink resource creation 
(#4266)
29242e316 is described below

commit 29242e3164befb5a837f1964e476ea61f6f058be
Author: woofyzhao <[email protected]>
AuthorDate: Thu May 26 10:52:16 2022 +0800

    [INLONG-4263][Manager] Support HBase sink resource creation (#4266)
---
 .../pojo/sink/hbase/HbaseColumnFamilyInfo.java     |  65 +++++++++
 .../common/pojo/sink/hbase/HbaseSinkDTO.java       |  13 ++
 .../common/pojo/sink/hbase/HbaseTableInfo.java     |  36 +++++
 .../service/resource/hbase/HbaseApiUtils.java      | 148 ++++++++++++++++++++
 .../resource/hbase/HbaseResourceOperator.java      | 152 +++++++++++++++++++++
 .../resource/iceberg/IcebergCatalogUtils.java      |   2 +-
 .../manager/service/sort/util/SinkInfoUtils.java   |  54 ++++++--
 .../inlong/sort/configuration/Constants.java       |   4 +
 .../inlong/sort/protocol/sink/HbaseSinkInfo.java   | 110 +++++++++++++++
 .../apache/inlong/sort/protocol/sink/SinkInfo.java |  15 +-
 10 files changed, 579 insertions(+), 20 deletions(-)

diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseColumnFamilyInfo.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseColumnFamilyInfo.java
new file mode 100644
index 000000000..f3f54afe5
--- /dev/null
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseColumnFamilyInfo.java
@@ -0,0 +1,65 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.manager.common.pojo.sink.hbase;
+
+import com.fasterxml.jackson.databind.DeserializationFeature;
+import com.fasterxml.jackson.databind.ObjectMapper;
+import io.swagger.annotations.ApiModelProperty;
+import lombok.AllArgsConstructor;
+import lombok.Builder;
+import lombok.Data;
+import lombok.NoArgsConstructor;
+import org.apache.commons.lang3.StringUtils;
+import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
+import org.apache.inlong.manager.common.exceptions.BusinessException;
+
+import javax.validation.constraints.NotNull;
+
+/**
+ * Hbase column family info
+ */
+@Data
+@Builder
+@NoArgsConstructor
+@AllArgsConstructor
+public class HbaseColumnFamilyInfo {
+
+    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();
+
+    @ApiModelProperty("Column family name")
+    private String cfName;
+
+    @ApiModelProperty("Column family ttl")
+    private Integer ttl;
+
+    /**
+     * Get the extra param from the Json
+     */
+    public static HbaseColumnFamilyInfo getFromJson(@NotNull String extParams) 
{
+        if (StringUtils.isEmpty(extParams)) {
+            return new HbaseColumnFamilyInfo();
+        }
+        try {
+            
OBJECT_MAPPER.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, 
false);
+            return OBJECT_MAPPER.readValue(extParams, 
HbaseColumnFamilyInfo.class);
+        } catch (Exception e) {
+            throw new 
BusinessException(ErrorCodeEnum.SINK_INFO_INCORRECT.getMessage());
+        }
+    }
+
+}
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseSinkDTO.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseSinkDTO.java
index d246f5ea2..8b2499bbc 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseSinkDTO.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseSinkDTO.java
@@ -28,6 +28,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 
 import javax.validation.constraints.NotNull;
+import java.util.List;
 import java.util.Map;
 
 /**
@@ -97,4 +98,16 @@ public class HbaseSinkDTO {
         }
     }
 
+    /**
+     * Get hbase table info
+     */
+    public static HbaseTableInfo getHbaseTableInfo(HbaseSinkDTO hbaseInfo, 
List<HbaseColumnFamilyInfo> columnFamilies) {
+        HbaseTableInfo info = new HbaseTableInfo();
+        info.setNamespace(hbaseInfo.getNamespace());
+        info.setTableName(hbaseInfo.getTableName());
+        info.setTblProperties(hbaseInfo.getProperties());
+        info.setColumnFamilies(columnFamilies);
+        return info;
+    }
+
 }
diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseTableInfo.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseTableInfo.java
new file mode 100644
index 000000000..a0e875a05
--- /dev/null
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/pojo/sink/hbase/HbaseTableInfo.java
@@ -0,0 +1,36 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.manager.common.pojo.sink.hbase;
+
+import lombok.Data;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * Hbase table info
+ */
+@Data
+public class HbaseTableInfo {
+
+    private String namespace;
+    private String tableName;
+    private String tableDesc;
+    private Map<String, Object> tblProperties;
+    private List<HbaseColumnFamilyInfo> columnFamilies;
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/hbase/HbaseApiUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/hbase/HbaseApiUtils.java
new file mode 100644
index 000000000..38ff865dd
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/hbase/HbaseApiUtils.java
@@ -0,0 +1,148 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.manager.service.resource.hbase;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.hbase.HBaseConfiguration;
+import org.apache.hadoop.hbase.NamespaceDescriptor;
+import org.apache.hadoop.hbase.TableName;
+import org.apache.hadoop.hbase.client.Admin;
+import org.apache.hadoop.hbase.client.ColumnFamilyDescriptor;
+import org.apache.hadoop.hbase.client.ColumnFamilyDescriptorBuilder;
+import org.apache.hadoop.hbase.client.Connection;
+import org.apache.hadoop.hbase.client.ConnectionFactory;
+import org.apache.hadoop.hbase.client.Table;
+import org.apache.hadoop.hbase.client.TableDescriptorBuilder;
+import org.apache.inlong.manager.common.pojo.sink.hbase.HbaseColumnFamilyInfo;
+import org.apache.inlong.manager.common.pojo.sink.hbase.HbaseTableInfo;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+
+/**
+ * Utils for hbase api
+ */
+public class HbaseApiUtils {
+
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(HbaseApiUtils.class);
+
+    private static final String HBASE_CONF_ZK_QUORUM = 
"hbase.zookeeper.quorum";
+    private static final String HBASE_CONF_ZNODE_PARENT = 
"zookeeper.znode.parent";
+
+    /**
+     * Get and verify hbase connection
+     */
+    private static Connection getConnection(String zkAddress, String zkNode) 
throws Exception {
+        Configuration config = HBaseConfiguration.create();
+
+        // ip1:port,ip2:port,...
+        config.set(HBASE_CONF_ZK_QUORUM, zkAddress);
+        config.set(HBASE_CONF_ZNODE_PARENT, zkNode);
+
+        return ConnectionFactory.createConnection(config);
+    }
+
+    /**
+     * Create hbase namespace
+     */
+    public static void createNamespace(String zkAddress, String zkNode, String 
namespace) throws Exception {
+        if (namespace == null || namespace.isEmpty()) {
+            return;
+        }
+        try (Connection conn = getConnection(zkAddress, zkNode)) {
+            Admin admin = conn.getAdmin();
+            if (Arrays.asList(admin.listNamespaces()).contains(namespace)) {
+                LOGGER.info("hbase namespace {} already exists", namespace);
+                return;
+            }
+            
admin.createNamespace(NamespaceDescriptor.create(namespace).build());
+            LOGGER.info("hbase namespace {} created", namespace);
+        }
+    }
+
+    /**
+     * Create hbase table
+     */
+    public static void createTable(String zkAddress, String zkNode, 
HbaseTableInfo tableInfo) throws Exception {
+        TableName tableName = TableName.valueOf(tableInfo.getNamespace(), 
tableInfo.getTableName());
+        TableDescriptorBuilder desc = 
TableDescriptorBuilder.newBuilder(tableName);
+        for (HbaseColumnFamilyInfo cf : tableInfo.getColumnFamilies()) {
+            // properties of column families can also be set here with builder
+            // frontend doesn't introduce much at this moment
+            
desc.setColumnFamily(ColumnFamilyDescriptorBuilder.of(cf.getCfName()));
+        }
+        try (Connection conn = getConnection(zkAddress, zkNode)) {
+            Admin admin = conn.getAdmin();
+            admin.createTable(desc.build());
+        }
+    }
+
+    /**
+     * Check hbase table already exists or not
+     */
+    public static boolean tableExists(String zkAddress, String zkNode, String 
namespace, String qualifier)
+            throws Exception {
+        TableName tableName = TableName.valueOf(namespace, qualifier);
+        try (Connection conn = getConnection(zkAddress, zkNode)) {
+            Admin admin = conn.getAdmin();
+            return admin.tableExists(tableName);
+        }
+    }
+
+    /**
+     * Query hbase table column families
+     */
+    public static List<HbaseColumnFamilyInfo> getColumnFamilies(String 
zkAddress, String zkNode, String namespace,
+            String qualifier) throws Exception {
+        List<HbaseColumnFamilyInfo> cfList = new ArrayList<>();
+        TableName tableName = TableName.valueOf(namespace, qualifier);
+        try (Connection conn = getConnection(zkAddress, zkNode)) {
+            Table table = conn.getTable(tableName);
+            for (ColumnFamilyDescriptor cf : 
table.getDescriptor().getColumnFamilies()) {
+                HbaseColumnFamilyInfo info = new HbaseColumnFamilyInfo();
+                info.setCfName(cf.getNameAsString());
+                info.setTtl(cf.getTimeToLive());
+                cfList.add(info);
+            }
+        }
+        return cfList;
+    }
+
+    /**
+     * Add column families for hbase table
+     */
+    public static void addColumnFamilies(String zkAddress, String zkNode, 
String namespace, String qualifier,
+            List<HbaseColumnFamilyInfo> columnFamilies) throws Exception {
+        TableName tableName = TableName.valueOf(namespace, qualifier);
+        try (Connection conn = getConnection(zkAddress, zkNode)) {
+            Admin admin = conn.getAdmin();
+            admin.disableTable(tableName);
+            try {
+                for (HbaseColumnFamilyInfo info : columnFamilies) {
+                    admin.addColumnFamily(tableName, 
ColumnFamilyDescriptorBuilder.of(info.getCfName()));
+                }
+            } finally {
+                admin.enableTable(tableName);
+            }
+        }
+    }
+
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/hbase/HbaseResourceOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/hbase/HbaseResourceOperator.java
new file mode 100644
index 000000000..6264459d4
--- /dev/null
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/hbase/HbaseResourceOperator.java
@@ -0,0 +1,152 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.manager.service.resource.hbase;
+
+import org.apache.commons.collections.CollectionUtils;
+import org.apache.inlong.manager.common.enums.GlobalConstants;
+import org.apache.inlong.manager.common.enums.SinkStatus;
+import org.apache.inlong.manager.common.enums.SinkType;
+import org.apache.inlong.manager.common.exceptions.WorkflowException;
+import org.apache.inlong.manager.common.pojo.sink.SinkInfo;
+import org.apache.inlong.manager.common.pojo.sink.hbase.HbaseColumnFamilyInfo;
+import org.apache.inlong.manager.common.pojo.sink.hbase.HbaseSinkDTO;
+import org.apache.inlong.manager.common.pojo.sink.hbase.HbaseTableInfo;
+import org.apache.inlong.manager.dao.entity.StreamSinkFieldEntity;
+import org.apache.inlong.manager.dao.mapper.StreamSinkFieldEntityMapper;
+import org.apache.inlong.manager.service.resource.SinkResourceOperator;
+import org.apache.inlong.manager.service.sink.StreamSinkService;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.stereotype.Service;
+
+import java.util.ArrayList;
+import java.util.Comparator;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Set;
+
+import static java.util.stream.Collectors.toList;
+
+/**
+ * hbase resource operator
+ */
+@Service
+public class HbaseResourceOperator implements SinkResourceOperator {
+
+    private static final Logger LOGGER = 
LoggerFactory.getLogger(HbaseResourceOperator.class);
+
+    @Autowired
+    private StreamSinkService sinkService;
+    @Autowired
+    private StreamSinkFieldEntityMapper sinkFieldMapper;
+
+    @Override
+    public Boolean accept(SinkType sinkType) {
+        return SinkType.HBASE == sinkType;
+    }
+
+    /**
+     * Create hbase table according to the sink config
+     */
+    public void createSinkResource(SinkInfo sinkInfo) {
+        if (sinkInfo == null) {
+            LOGGER.warn("sink info was null, skip to create resource");
+            return;
+        }
+
+        if 
(SinkStatus.CONFIG_SUCCESSFUL.getCode().equals(sinkInfo.getStatus())) {
+            LOGGER.warn("sink resource [" + sinkInfo.getId() + "] already 
success, skip to create");
+            return;
+        } else if 
(GlobalConstants.DISABLE_CREATE_RESOURCE.equals(sinkInfo.getEnableCreateResource()))
 {
+            LOGGER.warn("create resource was disabled, skip to create for [" + 
sinkInfo.getId() + "]");
+            return;
+        }
+
+        this.createTable(sinkInfo);
+    }
+
+    private void createTable(SinkInfo sinkInfo) {
+        LOGGER.info("begin to create hbase table for sinkInfo={}", sinkInfo);
+
+        // Get all info from config
+        HbaseSinkDTO hbaseInfo = 
HbaseSinkDTO.getFromJson(sinkInfo.getExtParams());
+        List<HbaseColumnFamilyInfo> columnFamilies = 
getColumnFamilies(sinkInfo);
+        if (CollectionUtils.isEmpty(columnFamilies)) {
+            throw new IllegalArgumentException("no hbase column families 
specified");
+        }
+        HbaseTableInfo tableInfo = HbaseSinkDTO.getHbaseTableInfo(hbaseInfo, 
columnFamilies);
+
+        String zkAddress = hbaseInfo.getZookeeperQuorum();
+        String zkNode = hbaseInfo.getZookeeperZnodeParent();
+        String namespace = hbaseInfo.getNamespace();
+        String tableName = hbaseInfo.getTableName();
+
+        try {
+            // 1. create database if not exists
+            HbaseApiUtils.createNamespace(zkAddress, zkNode, namespace);
+
+            // 2. check if the table exists
+            boolean tableExists = HbaseApiUtils.tableExists(zkAddress, zkNode, 
namespace, tableName);
+
+            if (!tableExists) {
+                // 3. create table
+                HbaseApiUtils.createTable(zkAddress, zkNode, tableInfo);
+            } else {
+                // 4. or update table columns
+                List<HbaseColumnFamilyInfo> existColumnFamilies = 
HbaseApiUtils.getColumnFamilies(zkAddress, zkNode,
+                                namespace, tableName).stream()
+                        
.sorted(Comparator.comparing(HbaseColumnFamilyInfo::getCfName)).collect(toList());
+                List<HbaseColumnFamilyInfo> requestColumnFamilies = 
tableInfo.getColumnFamilies().stream()
+                        
.sorted(Comparator.comparing(HbaseColumnFamilyInfo::getCfName)).collect(toList());
+                List<HbaseColumnFamilyInfo> newColumnFamilies = 
requestColumnFamilies.stream()
+                        .skip(existColumnFamilies.size()).collect(toList());
+
+                if (CollectionUtils.isNotEmpty(newColumnFamilies)) {
+                    HbaseApiUtils.addColumnFamilies(zkAddress, zkNode, 
namespace, tableName, newColumnFamilies);
+                    LOGGER.info("{} column families added for table {}", 
newColumnFamilies.size(), tableName);
+                }
+            }
+            String info = "success to create hbase resource";
+            sinkService.updateStatus(sinkInfo.getId(), 
SinkStatus.CONFIG_SUCCESSFUL.getCode(), info);
+            LOGGER.info(info + " for sinkInfo = {}", info);
+        } catch (Throwable e) {
+            String errMsg = "create hbase table failed: " + e.getMessage();
+            LOGGER.error(errMsg, e);
+            sinkService.updateStatus(sinkInfo.getId(), 
SinkStatus.CONFIG_FAILED.getCode(), errMsg);
+            throw new WorkflowException(errMsg);
+        }
+    }
+
+    private List<HbaseColumnFamilyInfo> getColumnFamilies(SinkInfo sinkInfo) {
+        List<StreamSinkFieldEntity> fieldList = 
sinkFieldMapper.selectBySinkId(sinkInfo.getId());
+        Set<String> seen = new HashSet<>();
+
+        List<HbaseColumnFamilyInfo> columnFamilies = new ArrayList<>();
+        for (StreamSinkFieldEntity field : fieldList) {
+            HbaseColumnFamilyInfo columnFamily = 
HbaseColumnFamilyInfo.getFromJson(field.getExtrParam());
+            if (seen.contains(columnFamily.getCfName())) {
+                continue;
+            }
+            seen.add(columnFamily.getCfName());
+            columnFamilies.add(columnFamily);
+        }
+
+        return columnFamilies;
+    }
+}
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/iceberg/IcebergCatalogUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/iceberg/IcebergCatalogUtils.java
index d73bfb28f..1d44c09b3 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/iceberg/IcebergCatalogUtils.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/iceberg/IcebergCatalogUtils.java
@@ -226,7 +226,7 @@ public class IcebergCatalogUtils {
 
     /**
      * Update iceberg table column schema.
-     * It's unfortunate that the updating api is different from the creating 
api so the column type switch is
+     * It's unfortunate that the updating api is different from the creating 
api so the partition type switch is
      * repeated here.
      */
     private static void updateColumnSpec(IcebergColumnInfo column, 
UpdatePartitionSpec builder) {
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SinkInfoUtils.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SinkInfoUtils.java
index 7a683b4b1..2d17d8bd4 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SinkInfoUtils.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/sort/util/SinkInfoUtils.java
@@ -28,6 +28,7 @@ import 
org.apache.inlong.manager.common.pojo.sink.SinkFieldBase;
 import org.apache.inlong.manager.common.pojo.sink.SinkResponse;
 import org.apache.inlong.manager.common.pojo.sink.ck.ClickHouseSinkResponse;
 import org.apache.inlong.manager.common.pojo.sink.es.ElasticsearchSinkResponse;
+import org.apache.inlong.manager.common.pojo.sink.hbase.HbaseSinkResponse;
 import org.apache.inlong.manager.common.pojo.sink.hive.HivePartitionField;
 import org.apache.inlong.manager.common.pojo.sink.hive.HiveSinkResponse;
 import org.apache.inlong.manager.common.pojo.sink.iceberg.IcebergSinkResponse;
@@ -38,6 +39,7 @@ import 
org.apache.inlong.sort.protocol.serialization.SerializationInfo;
 import org.apache.inlong.sort.protocol.sink.ClickHouseSinkInfo;
 import 
org.apache.inlong.sort.protocol.sink.ClickHouseSinkInfo.PartitionStrategy;
 import org.apache.inlong.sort.protocol.sink.ElasticsearchSinkInfo;
+import org.apache.inlong.sort.protocol.sink.HbaseSinkInfo;
 import org.apache.inlong.sort.protocol.sink.HiveSinkInfo;
 import 
org.apache.inlong.sort.protocol.sink.HiveSinkInfo.HiveFieldPartitionInfo;
 import org.apache.inlong.sort.protocol.sink.HiveSinkInfo.HiveFileFormat;
@@ -69,18 +71,27 @@ public class SinkInfoUtils {
             List<FieldInfo> sinkFields) {
         String sinkType = sinkResponse.getSinkType();
         SinkInfo sinkInfo;
-        if (SinkType.forType(sinkType) == SinkType.HIVE) {
-            sinkInfo = createHiveSinkInfo((HiveSinkResponse) sinkResponse, 
sinkFields);
-        } else if (SinkType.forType(sinkType) == SinkType.KAFKA) {
-            sinkInfo = createKafkaSinkInfo(sourceResponse, (KafkaSinkResponse) 
sinkResponse, sinkFields);
-        } else if (SinkType.SINK_ICEBERG.equals(sinkType)) {
-            sinkInfo = createIcebergSinkInfo((IcebergSinkResponse) 
sinkResponse, sinkFields);
-        } else if (SinkType.forType(sinkType) == SinkType.CLICKHOUSE) {
-            sinkInfo = createClickhouseSinkInfo((ClickHouseSinkResponse) 
sinkResponse, sinkFields);
-        } else if (SinkType.forType(sinkType) == SinkType.CLICKHOUSE) {
-            sinkInfo = createElasticsearchSinkInfo((ElasticsearchSinkResponse) 
sinkResponse, sinkFields);
-        } else {
-            throw new BusinessException(String.format("Unsupported SinkType 
{%s}", sinkType));
+        switch (SinkType.forType(sinkType)) {
+            case HIVE:
+                sinkInfo = createHiveSinkInfo((HiveSinkResponse) sinkResponse, 
sinkFields);
+                break;
+            case KAFKA:
+                sinkInfo = createKafkaSinkInfo(sourceResponse, 
(KafkaSinkResponse) sinkResponse, sinkFields);
+                break;
+            case ICEBERG:
+                sinkInfo = createIcebergSinkInfo((IcebergSinkResponse) 
sinkResponse, sinkFields);
+                break;
+            case CLICKHOUSE:
+                sinkInfo = createClickhouseSinkInfo((ClickHouseSinkResponse) 
sinkResponse, sinkFields);
+                break;
+            case HBASE:
+                sinkInfo = createHbaseSinkInfo((HbaseSinkResponse) 
sinkResponse, sinkFields);
+                break;
+            case ELASTICSEARCH:
+                sinkInfo = 
createElasticsearchSinkInfo((ElasticsearchSinkResponse) sinkResponse, 
sinkFields);
+                break;
+            default:
+                throw new BusinessException(String.format("Unsupported 
SinkType {%s}", sinkType));
         }
         return sinkInfo;
     }
@@ -237,6 +248,25 @@ public class SinkInfoUtils {
         }
     }
 
+    /**
+     * Creat HBase sink info.
+     */
+    private static HbaseSinkInfo createHbaseSinkInfo(HbaseSinkResponse 
sinkResponse, List<FieldInfo> sinkFields) {
+        if (StringUtils.isEmpty(sinkResponse.getZookeeperQuorum())) {
+            throw new BusinessException(String.format("HBase={%s} zookeeper 
quorum url cannot be empty", sinkResponse));
+        } else if 
(StringUtils.isEmpty(sinkResponse.getZookeeperZnodeParent())) {
+            throw new BusinessException(String.format("HBase={%s} zookeeper 
node cannot be empty", sinkResponse));
+        } else if (StringUtils.isEmpty(sinkResponse.getTableName())) {
+            throw new BusinessException(String.format("HBase={%s} table name 
cannot be empty", sinkResponse));
+        }
+
+        return new HbaseSinkInfo(sinkFields.toArray(new FieldInfo[0]), 
sinkResponse.getZookeeperQuorum(),
+                sinkResponse.getZookeeperZnodeParent(), 
sinkResponse.getNamespace(), sinkResponse.getTableName(),
+                sinkResponse.getSinkBufferFlushMaxSize(), 
sinkResponse.getSinkBufferFlushMaxSize(),
+                sinkResponse.getSinkBufferFlushInterval());
+
+    }
+
     /**
      * Creat Elasticsearch sink info.
      */
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/configuration/Constants.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/configuration/Constants.java
index 6bd0138e1..e5f836f34 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/configuration/Constants.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/configuration/Constants.java
@@ -44,6 +44,10 @@ public class Constants {
 
     public static final String SINK_TYPE_KAFKA = "kafka";
 
+    public static final String SINK_TYPE_HBASE = "hbase";
+
+    public static final String SINK_TYPE_ES = "elasticsearch";
+
     public static final String METRIC_DATA_OUTPUT_TAG_ID = 
"metric_data_side_output";
 
     public static final int METRIC_AUDIT_ID_FOR_INPUT = 7;
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/sink/HbaseSinkInfo.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/sink/HbaseSinkInfo.java
new file mode 100644
index 000000000..97aca92c8
--- /dev/null
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/sink/HbaseSinkInfo.java
@@ -0,0 +1,110 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *     http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.inlong.sort.protocol.sink;
+
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonCreator;
+import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
+import org.apache.inlong.sort.protocol.FieldInfo;
+
+import javax.annotation.Nullable;
+
+/**
+ * hbase sink resource info
+ */
+public class HbaseSinkInfo extends SinkInfo {
+
+    private static final long serialVersionUID = -7651732809476005186L;
+
+    @JsonProperty("zk_address")
+    private final String zkAddress;
+
+    @JsonProperty("zk_node")
+    private final String zkNode;
+
+    @JsonProperty("namespace")
+    private final String namespace;
+
+    @JsonProperty("table")
+    private final String tableName;
+
+    @JsonProperty("flush_max_size")
+    private final String flushMaxSize;
+
+    @JsonProperty("flush_max_rows")
+    private final String flushMaxRows;
+
+    @JsonProperty("flush_interval")
+    private final String flushInterval;
+
+    @JsonCreator
+    public HbaseSinkInfo(
+            @JsonProperty("fields") FieldInfo[] fields,
+            @JsonProperty("zk_address") String zkAddress,
+            @JsonProperty("zk_node") String zkNode,
+            @JsonProperty("namespace") @Nullable String namespace,
+            @JsonProperty("table") String tableName,
+            @JsonProperty("flush_max_size") @Nullable String flushMaxSize,
+            @JsonProperty("flush_max_rows") @Nullable String flushMaxRows,
+            @JsonProperty("flush_interval") @Nullable String flushInterval) {
+        super(fields);
+        this.zkAddress = zkAddress;
+        this.zkNode = zkNode;
+        this.namespace = namespace;
+        this.tableName = tableName;
+        this.flushMaxSize = flushMaxSize;
+        this.flushMaxRows = flushMaxRows;
+        this.flushInterval = flushInterval;
+    }
+
+    @JsonProperty("zk_address")
+    public String getZkAddress() {
+        return zkAddress;
+    }
+
+    @JsonProperty("zk_node")
+    public String getZkNode() {
+        return zkNode;
+    }
+
+    @JsonProperty("namespace")
+    public String getNamespace() {
+        return namespace;
+    }
+
+    @JsonProperty("table")
+    public String getTableName() {
+        return tableName;
+    }
+
+    @JsonProperty("flush_max_size")
+    public String getFlushMaxSize() {
+        return flushMaxSize;
+    }
+
+    @JsonProperty("flush_max_rows")
+    public String getFlushMaxRows() {
+        return flushMaxRows;
+    }
+
+    @JsonProperty("flush_interval")
+    public String getFlushInterval() {
+        return flushInterval;
+    }
+
+}
diff --git 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/sink/SinkInfo.java
 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/sink/SinkInfo.java
index 309fa74a0..2ee0a21ec 100644
--- 
a/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/sink/SinkInfo.java
+++ 
b/inlong-sort/sort-common/src/main/java/org/apache/inlong/sort/protocol/sink/SinkInfo.java
@@ -17,11 +17,6 @@
 
 package org.apache.inlong.sort.protocol.sink;
 
-import static com.google.common.base.Preconditions.checkNotNull;
-
-import java.io.Serializable;
-import java.util.Arrays;
-
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonProperty;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonSubTypes;
 import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonSubTypes.Type;
@@ -29,6 +24,11 @@ import 
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.annotation.JsonTyp
 import org.apache.inlong.sort.configuration.Constants;
 import org.apache.inlong.sort.protocol.FieldInfo;
 
+import java.io.Serializable;
+import java.util.Arrays;
+
+import static com.google.common.base.Preconditions.checkNotNull;
+
 /**
  * The base class of the data sink in the metadata.
  */
@@ -40,8 +40,9 @@ import org.apache.inlong.sort.protocol.FieldInfo;
         @Type(value = ClickHouseSinkInfo.class, name = 
Constants.SINK_TYPE_CLICKHOUSE),
         @Type(value = HiveSinkInfo.class, name = Constants.SINK_TYPE_HIVE),
         @Type(value = KafkaSinkInfo.class, name = Constants.SINK_TYPE_KAFKA),
-        @Type(value = HiveSinkInfo.class, name = Constants.SINK_TYPE_HIVE),
-        @Type(value = IcebergSinkInfo.class, name = 
Constants.SINK_TYPE_ICEBERG)}
+        @Type(value = IcebergSinkInfo.class, name = 
Constants.SINK_TYPE_ICEBERG),
+        @Type(value = HbaseSinkInfo.class, name = Constants.SINK_TYPE_HBASE),
+        @Type(value = ElasticsearchSinkInfo.class, name = 
Constants.SINK_TYPE_ES)}
 )
 public abstract class SinkInfo implements Serializable {
 

Reply via email to