hbase操作,创建命名空间,创建hbase表,添加数据,删除数据,通过RowKey获取数据,扫描数据,获取所有的表,获取ColumnFamily,删除表,获取hbase命名空间
·
1、创建hbase工程

2 编写config.properties
# hbase信息
hbase.zookeeper.quorum=xxx.xxx.xxx.xxx,xxx.xxx.xxx.xxx,xxx.xxx.xxx.xxx
hbase.zookeeper.property.clientPort=2181
3 编写Config
package com.test.hbasecopy.beans;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.PropertySource;
/**
* @author tuzuoquan
* @version 1.0
* @ClassName Config
* @description TODO
* @date 2020/8/28 10:28
**/
@Configuration
@PropertySource({
"classpath:config.properties"
})
public class Config {
/** 要迁移数据的源hbase数据信息 **/
@Value("${hbase.zookeeper.quorum}")
public String hbaseZkQuorum;
@Value("${hbase.zookeeper.property.clientPort}")
public String hbaseZkClientPort;
public String getHbaseZkQuorum() {
return hbaseZkQuorum;
}
public void setHbaseZkQuorum(String hbaseZkQuorum) {
this.hbaseZkQuorum = hbaseZkQuorum;
}
public String getHbaseZkClientPort() {
return hbaseZkClientPort;
}
public void setHbaseZkClientPort(String hbaseZkClientPort) {
this.hbaseZkClientPort = hbaseZkClientPort;
}
}
4 BooleanDataPair
package com.test.hbasecopy.domain.base;
import com.test.hbasecopy.domain.DataPair;
public class BooleanDataPair extends DataPair {
private static final long serialVersionUID = 1L;
private boolean value;
public BooleanDataPair() {
}
public BooleanDataPair(String columnName, boolean value) {
this.columnName = columnName;
this.value = value;
}
public boolean getValue() {
return value;
}
public void setValue(boolean value) {
this.value = value;
}
}
5 DoubleDataPair
package com.test.hbasecopy.domain.base;
import com.test.hbasecopy.domain.DataPair;
public class DoubleDataPair extends DataPair {
private static final long serialVersionUID = 1L;
private double value;
public DoubleDataPair() {
}
public DoubleDataPair(String columnName, double value) {
this.columnName = columnName;
this.value = value;
}
public double getValue() {
return value;
}
public void setValue(double value) {
this.value = value;
}
}
6 FloatDataPair
package com.test.hbasecopy.domain.base;
import com.test.hbasecopy.domain.DataPair;
public class FloatDataPair extends DataPair {
private static final long serialVersionUID = 1L;
private float value;
public FloatDataPair() {
}
public FloatDataPair(String columnName, float value) {
this.columnName = columnName;
this.value = value;
}
public float getValue() {
return value;
}
public void setValue(float value) {
this.value = value;
}
}
7 IntDataPair
package com.test.hbasecopy.domain.base;
import com.test.hbasecopy.domain.DataPair;
public class IntDataPair extends DataPair {
private static final long serialVersionUID = 1L;
private int value;
public IntDataPair() {
}
public IntDataPair(String columnName, int value) {
this.columnName = columnName;
this.value = value;
}
public int getValue() {
return value;
}
public void setValue(int value) {
this.value = value;
}
}
8 LongDataPair
package com.test.hbasecopy.domain.base;
import com.test.hbasecopy.domain.DataPair;
public class LongDataPair extends DataPair {
private static final long serialVersionUID = 1L;
private long value;
public LongDataPair() {
}
public LongDataPair(String columnName, long value) {
this.columnName = columnName;
this.value = value;
}
public long getValue() {
return value;
}
public void setValue(long value) {
this.value = value;
}
}
9 ShortDataPair
package com.test.hbasecopy.domain.base;
import com.test.hbasecopy.domain.DataPair;
public class ShortDataPair extends DataPair {
private static final long serialVersionUID = 1L;
private short value;
public ShortDataPair() {
}
public ShortDataPair(String columnName, short value) {
this.columnName = columnName;
this.value = value;
}
public short getValue() {
return value;
}
public void setValue(short value) {
this.value = value;
}
}
10 StringDataPair
package com.test.hbasecopy.domain.base;
import com.test.hbasecopy.domain.DataPair;
public class StringDataPair extends DataPair {
private static final long serialVersionUID = 1L;
private String value;
public StringDataPair() {
}
public StringDataPair(String columnName, String value) {
this.columnName = columnName;
this.value = value;
}
public String getValue() {
return value;
}
public void setValue(String value) {
this.value = value;
}
}
11 ColumnFamilyData
package com.test.hbasecopy.domain;
import java.io.Serializable;
import java.util.ArrayList;
import java.util.List;
public class ColumnFamilyData implements Serializable {
private static final long serialVersionUID = 1L;
private String columnFamily;
private List<DataPair> dataPairs;
public static ColumnFamilyData builder() {
return new ColumnFamilyData();
}
public String getColumnFamily() {
return columnFamily;
}
public ColumnFamilyData setColumnFamily(String columnFamily) {
this.columnFamily = columnFamily;
return this;
}
public List<DataPair> getDataPairs() {
return dataPairs;
}
public ColumnFamilyData setDataPairs(List<DataPair> dataPairs) {
this.dataPairs = dataPairs;
return this;
}
public ColumnFamilyData addColumnFamilyData(DataPair dataPair) {
if (this.dataPairs == null) {
dataPairs = new ArrayList<DataPair>();
}
dataPairs.add(dataPair);
return this;
}
}
12 DataPair
package com.test.hbasecopy.domain;
import java.io.Serializable;
public class DataPair implements Serializable {
private static final long serialVersionUID = 1L;
/**
* 列名:HBase列修饰符
*/
protected String columnName;
public String getColumnName() {
return columnName;
}
public void setColumnName(String columnName) {
this.columnName = columnName;
}
}
13 Deleter
package com.test.hbasecopy.domain;
import java.io.Serializable;
import java.util.ArrayList;
import java.util.List;
public class Deleter implements Serializable {
private static final long serialVersionUID = 1L;
private String rowKey;
private String tableName;
private List<FamilyQualifierMapper> familyQualifierMappers;
public static Deleter builder() {
return new Deleter();
}
public String getRowKey() {
return rowKey;
}
public Deleter setRowKey(String rowKey) {
this.rowKey = rowKey;
return this;
}
public String getTableName() {
return tableName;
}
public Deleter setTableName(String tableName) {
this.tableName = tableName;
return this;
}
public List<FamilyQualifierMapper> getFamilyQualifierMappers() {
return familyQualifierMappers;
}
public void setFamilyQualifierMappers(List<FamilyQualifierMapper> familyQualifierMappers) {
this.familyQualifierMappers = familyQualifierMappers;
}
public Deleter addFamilyQualifierMapper(FamilyQualifierMapper familyQualifierMapper) {
if (this.familyQualifierMappers == null) {
this.familyQualifierMappers = new ArrayList<FamilyQualifierMapper>();
}
familyQualifierMappers.add(familyQualifierMapper);
return this;
}
}
14 FamilyQualifierMapper
package com.test.hbasecopy.domain;
import java.io.Serializable;
public class FamilyQualifierMapper implements Serializable {
private static final long serialVersionUID = 1L;
private String family;
private String qualifier;
public FamilyQualifierMapper() {
}
public FamilyQualifierMapper(String family, String qualifier) {
this.family = family;
this.qualifier = qualifier;
}
public String getFamily() {
return family;
}
public void setFamily(String family) {
this.family = family;
}
public String getQualifier() {
return qualifier;
}
public void setQualifier(String qualifier) {
this.qualifier = qualifier;
}
}
15 RowData
package com.test.hbasecopy.domain;
import java.io.Serializable;
import java.util.ArrayList;
import java.util.List;
public class RowData implements Serializable {
private static final long serialVersionUID = 1L;
private String rowKey;
private String tableName;
private List<ColumnFamilyData> columnFamilyDatas;
public static RowData builder() {
return new RowData();
}
public String getRowKey() {
return rowKey;
}
public RowData setRowKey(String rowKey) {
this.rowKey = rowKey;
return this;
}
public String getTableName() {
return tableName;
}
public RowData setTableName(String tableName) {
this.tableName = tableName;
return this;
}
public List<ColumnFamilyData> getColumnFamilyDatas() {
return columnFamilyDatas;
}
public void setColumnFamilyDatas(List<ColumnFamilyData> columnFamilyDatas) {
this.columnFamilyDatas = columnFamilyDatas;
}
public RowData addColumnFamilyData(ColumnFamilyData columnFamilyData) {
if (this.columnFamilyDatas == null) {
columnFamilyDatas = new ArrayList<ColumnFamilyData>();
}
columnFamilyDatas.add(columnFamilyData);
return this;
}
}
16 Scanner
package com.test.hbasecopy.domain;
import java.io.Serializable;
import java.util.ArrayList;
import java.util.List;
public class Scanner implements Serializable {
private static final long serialVersionUID = 1L;
private List<FamilyQualifierMapper> familyQualifierMappers;
public static Scanner builder() {
return new Scanner();
}
public Scanner addFamilyQualifierMapper(FamilyQualifierMapper familyQualifierMapper) {
if (this.familyQualifierMappers == null) {
this.familyQualifierMappers = new ArrayList<FamilyQualifierMapper>();
}
familyQualifierMappers.add(familyQualifierMapper);
return this;
}
public List<FamilyQualifierMapper> getFamilyQualifierMappers() {
return this.familyQualifierMappers;
}
public void setFamilyQualifierMappers(List<FamilyQualifierMapper> familyQualifierMappers) {
this.familyQualifierMappers = familyQualifierMappers;
}
}
17 HbaseService
package com.test.hbasecopy.service;
import com.test.hbasecopy.domain.Deleter;
import com.test.hbasecopy.domain.RowData;
import com.test.hbasecopy.domain.Scanner;
import org.apache.hadoop.hbase.HTableDescriptor;
import org.apache.hadoop.hbase.NamespaceDescriptor;
import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.client.Connection;
import java.util.List;
import java.util.Map;
public interface HbaseService {
/**
* Hbase每条数据最大版本号
*/
final int maxVersion = 3;
/**
* 一次scan取到内存的记录数
*/
final int cacheSize = 100;
/**
* 生成Coonnection
* @param connection
*/
void setConnection(Connection connection);
/**
* 创建命名空间
*
* @param nameSpace
*/
void createNameSpace(String nameSpace);
/**
* 获取命名空间,如果不存在则创建.
*
* @param nameSpace
*/
void createHbaseNameSpaceIfNotExit(String nameSpace);
/**
* 获取所有命名空间
*
* @param nameSpace
* @return
*/
NamespaceDescriptor[] nameSpaceVerification(String nameSpace);
/**
* 创建表
*
* @param namespace
* @param tableName 表名
* @param families 列族名集合
*/
void createTable(String namespace, String tableName, List<String> families);
void addData(String namespace, RowData rowData);
void deleteData(String namespace, Deleter deleter);
Map<String, byte[]> getDataByRowKey(String namespace, String rowKey, String tableName, Scanner scanner);
Map<String, byte[]> getDataByRowKey(String namespace, String rowKey, String tableName);
List<Map<String, byte[]>> scanData(String namespace, String tableName, String startRowKey, String endRowKey,
Scanner scanner, int page, int pageSize);
List<Map<String, byte[]>> scanData(String namespace, String tableName, String startRowKey, String endRowKey,
int page, int pageSize);
/**
* HBase查询所有表
*
* @param namespace
* @return
*/
TableName[] queryAllTables(String namespace);
/**
* 获取列族信息
*
* @param tableName
* @return
*/
HTableDescriptor getColumnFamily(String tableName);
/**
* 删除表
*
* @param hbaseTables
* @return
*/
void deleteTables(TableName[] hbaseTables);
/**
* 检查命名空间是否存在.
* @param nameSpace
* @return
*/
boolean findHbaseNameSpace(String nameSpace);
}
18 HbaseServiceImpl
package com.tianque.hbasecopy.service.impl;
import com.test.hbasecopy.beans.Config;
import com.test.hbasecopy.domain.*;
import com.test.hbasecopy.domain.base.*;
import com.test.hbasecopy.service.HbaseService;
import com.test.hbasecopy.utils.HbaseUtil;
import org.apache.hadoop.hbase.*;
import org.apache.hadoop.hbase.client.*;
import org.apache.hadoop.hbase.util.Bytes;
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.HashMap;
import java.util.List;
import java.util.Map;
/**
* @author tuzuoquan
* @version 1.0
* @ClassName HbaseServiceImpl
* @description TODO
* @date 2020/8/28 11:29
**/
@Service(value = "hbaseService")
public class HbaseServiceImpl implements HbaseService {
private static final Logger log = LoggerFactory.getLogger(HbaseServiceImpl.class);
private Connection connection;
@Autowired
Config config;
// /** 构造方法:初始化时获取HBase的连接 */
// public HbaseServiceImpl() {
// try {
// connection = HbaseUtil.getHbaseConnection(config.getHbaseZkQuorum(),config.getHbaseZkClientPort());
// } catch (Exception e) {
// log.error("获取hbase连接失败", e);
// }
// }
/**
* @param connection
*/
@Override
public void setConnection(Connection connection) {
this.connection = connection;
}
/**
* 命名空间初始化
* @param nameSpace
*/
@Override
public void createNameSpace(String nameSpace) {
Admin admin = null;
try {
admin = connection.getAdmin();
admin.createNamespace(NamespaceDescriptor.create(nameSpace).build());
} catch (Exception e) {
throw new RuntimeException("创建Hbase命名空间失败:", e);
} finally {
try {
if (admin != null) {
admin.close();// 由于connection已经在admin中关闭,所以不必关闭connection
}
} catch (Exception e) {
log.error("关闭HAdmin异常:", e);
}
}
}
/**
* 获取命名空间,如果不存在则创建.
* @param nameSpace
*/
@Override
public void createHbaseNameSpaceIfNotExit(String nameSpace) {
Admin admin = null;
try{
//获取管理对象
admin = connection.getAdmin();
NamespaceDescriptor[] listNamespaceDescriptors = admin.listNamespaceDescriptors();
List<String> namespaceList = new ArrayList<String>();
for(int i=0;i<listNamespaceDescriptors.length;i++){
String name = listNamespaceDescriptors[i].getName();
namespaceList.add(name);
}
//如果命名空间不存在, 则创建命名空间.
if(!namespaceList.contains(nameSpace)){
createNameSpace(nameSpace);
}
}catch(Exception e){
throw new RuntimeException("获取Hbase命名空间失败:", e);
}
}
@Override
public NamespaceDescriptor[] nameSpaceVerification(String nameSpace) {
Admin admin = null;
try {
admin = connection.getAdmin();
NamespaceDescriptor[] namespaceDescriptor = admin.listNamespaceDescriptors();
return namespaceDescriptor;
} catch (Exception e) {
throw new RuntimeException("创建Hbase命名空间失败:", e);
} finally {
try {
if (admin != null) {
admin.close();// 由于connection已经在admin中关闭,所以不必关闭connection
}
} catch (Exception e) {
log.error("关闭HAdmin异常:", e);
}
}
}
@Override
public void createTable(String namespace, String tableName, List<String> families) {
Admin admin = null;
try {
admin = connection.getAdmin();
HTableDescriptor table = new HTableDescriptor(TableName.valueOf(namespace + ":" + tableName));
for (String columnFamily : families) {
table.addFamily(new HColumnDescriptor(columnFamily).setMaxVersions(maxVersion));
}
admin.createTable(table);
} catch (Exception e) {
throw new RuntimeException("创建Hbase表失败:", e);
} finally {
try {
if (admin != null) {
admin.close();// 由于connection已经在admin中关闭,所以不必关闭connection
}
} catch (Exception e) {
log.error("关闭HAdmin异常:", e);
}
}
}
@Override
public void addData(String namespace, RowData rowData) {
HTable table = null;
try {
Put put = new Put(Bytes.toBytes(rowData.getRowKey()));
table = (HTable) connection.getTable(TableName.valueOf(namespace + ":" + rowData.getTableName()));
for (ColumnFamilyData columnFamilyData : rowData.getColumnFamilyDatas()) {
for (DataPair dataPair : columnFamilyData.getDataPairs()) {
put.addColumn(Bytes.toBytes(columnFamilyData.getColumnFamily()),
Bytes.toBytes(dataPair.getColumnName()), getByte(dataPair));
}
}
table.put(put);
} catch (Exception e) {
throw new RuntimeException("新增数据到Hbase失败:", e);
} finally {
try {
if (table != null) {
table.close();
}
} catch (Exception e) {
log.error("关闭HTable异常:", e);
}
}
}
/**
*
* 删除某一列数据
*/
@Override
public void deleteData(String namespace, Deleter deleter) {
HTable table = null;
try {
Delete delete = new Delete(Bytes.toBytes(deleter.getRowKey()));
table = (HTable) connection.getTable(TableName.valueOf(namespace + ":" + deleter.getTableName()));
if (deleter.getFamilyQualifierMappers() != null) {
for (FamilyQualifierMapper familyQualifierMapper : deleter.getFamilyQualifierMappers()) {
delete.addColumns(Bytes.toBytes(familyQualifierMapper.getFamily()),
Bytes.toBytes(familyQualifierMapper.getQualifier()));
}
}
table.delete(delete);
} catch (Exception e) {
throw new RuntimeException("新增数据到Hbase失败:", e);
} finally {
try {
if (table != null) {
table.close();
}
} catch (Exception e) {
log.error("关闭HTable异常:", e);
}
}
}
@Override
public Map<String, byte[]> getDataByRowKey(String namespace, String rowKey, String tableName, Scanner scanner) {
HTable table = null;
Map<String, byte[]> result = new HashMap<String, byte[]>();
try {
table = (HTable) connection.getTable(TableName.valueOf(namespace + ":" + tableName));
Get get = new Get(Bytes.toBytes(rowKey));
if (scanner != null && scanner.getFamilyQualifierMappers() != null) {
for (FamilyQualifierMapper familyQualifierMapper : scanner.getFamilyQualifierMappers()) {
get.addColumn(Bytes.toBytes(familyQualifierMapper.getFamily()),
Bytes.toBytes(familyQualifierMapper.getQualifier()));
}
}
Result rs = table.get(get);
for (Cell cell : rs.rawCells()) {
result.put(Bytes.toString(CellUtil.cloneQualifier(cell)), CellUtil.cloneValue(cell));
}
} catch (Exception e) {
throw new RuntimeException("新增数据到Hbase失败:", e);
} finally {
try {
if (table != null) {
table.close();
}
} catch (Exception e) {
log.error("关闭HTable异常:", e);
}
}
return result;
}
@Override
public Map<String, byte[]> getDataByRowKey(String namespace, String rowKey, String tableName) {
return getDataByRowKey(namespace, rowKey, tableName, null);
}
@Override
public List<Map<String, byte[]>> scanData(String namespace, String tableName, String startRowKey, String endRowKey,
Scanner scanner, int page, int pageSize) {
HTable table = null;
List<Map<String, byte[]>> result = new ArrayList<Map<String, byte[]>>();
try {
table = (HTable) connection.getTable(TableName.valueOf(namespace + ":" + tableName));
Scan scan = new Scan();
scan.setCaching(cacheSize);
if (startRowKey != null) {
scan.setStartRow(Bytes.toBytes(startRowKey));
}
if (endRowKey != null) {
scan.setStopRow(Bytes.toBytes(endRowKey));
}
if (scanner != null && scanner.getFamilyQualifierMappers() != null) {
for (FamilyQualifierMapper familyQualifierMapper : scanner.getFamilyQualifierMappers()) {
scan.addColumn(Bytes.toBytes(familyQualifierMapper.getFamily()),
Bytes.toBytes(familyQualifierMapper.getQualifier()));
}
}
ResultScanner resultScanner = table.getScanner(scan);
int count = 0;
for (Result rs : resultScanner) {
count++;
if (count > (page - 1) * pageSize && count <= page * pageSize) {
Map<String, byte[]> resultMap = new HashMap<String, byte[]>();
String rowKey = null;
for (Cell cell : rs.rawCells()) {
resultMap.put(Bytes.toString(CellUtil.cloneQualifier(cell)), CellUtil.cloneValue(cell));
rowKey = cell.toString().split("/")[0];
}
resultMap.put("rowKey", rowKey.getBytes());
result.add(resultMap);
}
}
} catch (Exception e) {
throw new RuntimeException("新增数据到Hbase失败:", e);
} finally {
try {
if (table != null) {
table.close();
}
} catch (Exception e) {
log.error("关闭HTable异常:", e);
}
}
return result;
}
@Override
public List<Map<String, byte[]>> scanData(String namespace, String tableName, String startRowKey, String endRowKey,
int page, int pageSize) {
return scanData(namespace, tableName, startRowKey, endRowKey, null, page, pageSize);
}
private byte[] getByte(DataPair dataPair) {
byte[] value = null;
if (dataPair instanceof IntDataPair) {
value = Bytes.toBytes(((IntDataPair) dataPair).getValue());
} else if (dataPair instanceof ShortDataPair) {
value = Bytes.toBytes(((ShortDataPair) dataPair).getValue());
} else if (dataPair instanceof LongDataPair) {
value = Bytes.toBytes(((LongDataPair) dataPair).getValue());
} else if (dataPair instanceof DoubleDataPair) {
value = Bytes.toBytes(((DoubleDataPair) dataPair).getValue());
} else if (dataPair instanceof StringDataPair) {
value = Bytes.toBytes(((StringDataPair) dataPair).getValue());
} else if (dataPair instanceof BooleanDataPair) {
value = Bytes.toBytes(((BooleanDataPair) dataPair).getValue());
} else if (dataPair instanceof FloatDataPair) {
value = Bytes.toBytes(((FloatDataPair) dataPair).getValue());
}
return value;
}
/**
* HBase查询所有表
*/
@Override
public TableName[] queryAllTables(String namespace) {
Admin admin = null;
try {
admin = connection.getAdmin();
return admin.listTableNamesByNamespace(namespace);
} catch (Exception e) {
log.error("HBase查询所有表出错:", e);
} finally {
try {
if (admin != null) {
// 由于connection已经在admin中关闭,所以不必关闭connection
admin.close();
}
} catch (Exception e) {
log.error("关闭HAdmin异常:", e);
}
}
return null;
}
/**
* 获取列族信息
*/
@Override
public HTableDescriptor getColumnFamily(String tableName) {
Admin admin = null;
try {
admin = connection.getAdmin();
return admin.getTableDescriptor(TableName.valueOf(tableName));
} catch (Exception e) {
log.error("HBase查询所有表出错:", e);
} finally {
try {
if (admin != null) {
// 由于connection已经在admin中关闭,所以不必关闭connection
admin.close();
}
} catch (Exception e) {
log.error("关闭HAdmin异常:", e);
}
}
return null;
}
@Override
public void deleteTables(TableName[] hbaseTables) {
Admin admin = null;
try {
admin = connection.getAdmin();
for (int i = 0; i < hbaseTables.length; i++) {
TableName tableName = hbaseTables[i];
admin.disableTable(tableName);
admin.deleteTable(tableName);
}
} catch (Exception e) {
throw new RuntimeException("删除Hbase表失败:", e);
} finally {
try {
if (admin != null) {
// 由于connection已经在admin中关闭,所以不必关闭connection
admin.close();
}
} catch (Exception e) {
log.error("关闭HAdmin异常:", e);
}
}
}
@Override
public boolean findHbaseNameSpace(String nameSpace) {
Admin admin = null;
boolean exist = true;
boolean notExist = false;
try {
// 获取管理对象
admin = connection.getAdmin();
NamespaceDescriptor[] listNamespaceDescriptors = admin.listNamespaceDescriptors();
List<String> namespaceList = new ArrayList<String>();
for (int i = 0; i < listNamespaceDescriptors.length; i++) {
String name = listNamespaceDescriptors[i].getName();
namespaceList.add(name);
}
if (!namespaceList.contains(nameSpace)) {
return notExist;
}
} catch (Exception e) {
throw new RuntimeException("获取Hbase命名空间失败:", e);
}
return exist;
}
}
19 HbaseUtil
package com.test.hbasecopy.utils;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.hbase.HBaseConfiguration;
import org.apache.hadoop.hbase.client.Connection;
import org.apache.hadoop.hbase.client.ConnectionFactory;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* @author tuzuoquan
* @version 1.0
* @ClassName HbaseUtil
* @description TODO
* @date 2020/8/28 13:35
**/
public final class HbaseUtil {
private static final Logger logger = LoggerFactory.getLogger(HbaseUtil.class);
private static Connection connection = null;
/**
* 获取源hbase的Configuration
* @param hbaseZkQuorum
* @param hbaseZkClientPort
* @return
*/
public static Connection getHbaseConnection(String hbaseZkQuorum,String hbaseZkClientPort) {
if (null != connection) {
return connection;
}
try {
Configuration config = HBaseConfiguration.create();
config.set("hbase.zookeeper.quorum",hbaseZkQuorum);
config.set("hbase.zookeeper.property.clientPort",hbaseZkClientPort);
connection = ConnectionFactory.createConnection(config);
return connection;
} catch (Exception e) {
logger.error(">>>>>>>>>>>>获取源hbase Connection失败!");
throw new RuntimeException(e);
}
}
}
20 HbaseCopyApplication
package com.test.hbasecopy;
import org.springframework.beans.factory.annotation.Configurable;
import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.boot.autoconfigure.SpringBootApplication;
import org.springframework.context.annotation.ComponentScan;
@Configurable
@EnableAutoConfiguration
@ComponentScan
@SpringBootApplication
public class HbaseCopyApplication {
public static void main(String[] args) {
SpringApplication.run(HbaseCopyApplication.class, args);
System.out.println("开始执行.....");
}
}
21 application.properties
server.port=9090
server.servlet.context-path=/
22 HbaseDataCopyerController
package com.test.hbasecopy.controller;
import com.test.hbasecopy.beans.Config;
import com.test.hbasecopy.service.HbaseService;
import com.test.hbasecopy.utils.HbaseUtil;
import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.client.Connection;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.autoconfigure.EnableAutoConfiguration;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* @author tuzuoquan
* @version 1.0
* @ClassName HBaseDataCopyer
* @description TODO
* @date 2020/8/28 10:20
**/
@RestController
@EnableAutoConfiguration
@RequestMapping("/hbaseCopy")
public class HbaseDataCopyerController {
@Autowired
private Config config;
@Autowired
HbaseService hbaseService;
@RequestMapping("/test")
public String test() {
Connection connection = HbaseUtil.getHbaseConnection(config.getHbaseZkQuorum(),config.getHbaseZkClientPort());
hbaseService.setConnection(connection);
TableName[] tableNames = hbaseService.queryAllTables("default");
StringBuffer tables = new StringBuffer();
for(TableName tableName : tableNames) {
String name = tableName.getNameAsString();
System.out.println(name);
tables.append(name + "\r\n");
}
return tables.toString();
}
}
更多推荐


所有评论(0)