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 {