实践
自定义注解
@Target(ElementType.FIELD)
@Retention(RetentionPolicy.RUNTIME)
@Documented
public @interface HField {
String family() default "f1";
String qualifier() default "";
boolean isKey() default false;
}
使用注解
@HField(isKey = true)
private String key;
@HField(qualifier = "pf")
@Description("平台名")
private String platform;
JAVA类解析成Hbase的put
操作
Put put = HMapper.buildPut(g);
实现
@Slf4j
public final class HMapper {
private HMapper() {
JSON.DEFAULT_PARSER_FEATURE &= ~Feature.UseBigDecimal.getMask();
};
private static Map<Class<?>, Field> FIELD_CACHE = new ConcurrentHashMap<>();
private static String toUnderlineStyle(String name) {
String result = "";
for (char c : name.toCharArray()) {
if (c >= 'A' && c <= 'Z') {
c += 32;
if (result.length() > 0) {
result += '_';
}
}
result += c;
}
return result;
}
private static String toCamelStyle(String name) {
String result = "";
boolean underline = false;
for (char c : name.toCharArray()) {
if (c == '_') {
underline = true;
} else if (underline && c >= 'a' && c <= 'z') {
c -= 32;
underline = false;
result += c;
}
}
return result;
}
public static boolean mapper(Result r, Object t) {
if (null != r && !r.isEmpty()) {
for (Field field : ReflectUtil.listField(t.getClass(), HField.class)) {
try {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
if (hf.isKey()) {
field.set(t, Bytes.toString(r.getRow()));
} else {
byte[] family = Bytes.toBytes(hf.family());
List<byte[]> qualifiers = Lists.newArrayList();
if (!Strings.isNullOrEmpty(hf.qualifier())) {
qualifiers.add(Bytes.toBytes(hf.qualifier()));
}
qualifiers.add(Bytes.toBytes(field.getName()));
if(!toUnderlineStyle(field.getName()).equals(field.getName())) {
qualifiers.add(Bytes.toBytes(toUnderlineStyle(field.getName())));
}
if(!toCamelStyle(field.getName()).equals(field.getName())) {
qualifiers.add(Bytes.toBytes(toCamelStyle(field.getName())));
}
for (byte[] qualifier : qualifiers) {
if (r.containsColumn(family, qualifier)) {
byte[] value = r.getValue(family, qualifier);
if (field.getType() == byte[].class) {
field.set(t, value);
} else if (field.getType() == String.class) {
field.set(t, Bytes.toString(value));
} else if (field.getType() == Double.class) {
field.set(t, Bytes.toDouble(value));
} else if (field.getType() == Integer.class) {
try {
field.set(t, Bytes.toInt(value));
} catch (Exception e) {
}
} else if (field.getType() == Long.class) {
try {
field.set(t, Bytes.toLong(value));
} catch (Exception e) {
}
} else if (field.getType() == Date.class) {
if (null != field.getAnnotation(JsonFormat.class) && !Strings
.isNullOrEmpty(field.getAnnotation(JsonFormat.class).pattern())) {
field.set(t, TimeUtil.parse(Bytes.toString(value),
field.getAnnotation(JsonFormat.class).pattern()));
} else {
field.set(t, TimeUtil.parse(Bytes.toString(value)));
}
} else if (field.getType() == Boolean.class) {
field.set(t, Bytes.toBoolean(value));
} else if (field.getType() == Short.class) {
field.set(t, Bytes.toShort(value));
} else if (field.getType() == Float.class) {
field.set(t, Bytes.toFloat(value));
} else if (field.getType() == Character.class) {
field.set(t, Bytes.toInt(value));
} else if (field.getType().isArray()) {
JSONArray array = JSONArray.parseArray(Bytes.toString(value));
field.set(t, array.toJavaObject(field.getType()));
} else if (field.getType() == Set.class) {
JSONArray array = JSONArray.parseArray(Bytes.toString(value));
field.set(t, array.toJavaObject(field.getType()));
} else if (field.getType() == List.class) {
Object result = null;
String strValue = Bytes.toString(value);
JSONArray array = JSONArray.parseArray(strValue);
result = mapperActualType(field,strValue);
if (result == null){
result = array.toJavaObject(field.getType());
}
field.set(t,result);
} else if (field.getType() == Map.class) {
Map<String, Object> map = Maps.newLinkedHashMap();
JSONObject json = JSON.parseObject(Bytes.toString(value));
Type mapMainType = field.getGenericType();
if (mapMainType instanceof ParameterizedType) {
ParameterizedType parameterizedType = (ParameterizedType) mapMainType;
Type[] types = parameterizedType.getActualTypeArguments();
for (String key : json.keySet()) {
map.put(key, parse(types[1], json.get(key)));
}
}
field.set(t, map);
// field.set(t, json.toJavaObject(field.getType()));
} else if (field.getType() == Vector.class) {
JSONArray array = JSONArray.parseArray(Bytes.toString(value));
field.set(t, array.toJavaObject(field.getType()));
} else {
String str = Bytes.toString(value);
if (field.getType() == Object.class) {
field.set(t, JSON.parse(str,Feature.OrderedField));
} else {
field.set(t, JSON.toJavaObject(JSON.parseObject(str,Feature.OrderedField), field.getType()));
}
}
break;
}
}
}
}
} catch (Exception e) {
e.printStackTrace();
System.out.println(field.getName());
}
}
return true;
}
return false;
}
private static Object mapperActualType(Field field,String value) {
Class<?> fieldClazz = field.getType();
if(fieldClazz.isAssignableFrom(List.class)){
Type fc = field.getGenericType();
JSONArray jsonArray = JSON.parseArray(value);
if(fc instanceof ParameterizedType){
ParameterizedType pt = (ParameterizedType) fc;
Class genericClazz = (Class)pt.getActualTypeArguments()[0];
return jsonArray.toJavaList(genericClazz);
}
}
return null;
}
private static Object parse(Type type, Object obj) {
if (null != type && null != obj) {
if (type instanceof ParameterizedType) {
ParameterizedType parameterizedType = (ParameterizedType) type;
Type[] types = parameterizedType.getActualTypeArguments();
if (obj instanceof JSONArray && (((ParameterizedType) type).getRawType() == List.class
|| ((ParameterizedType) type).getRawType() == Set.class
|| ((ParameterizedType) type).getRawType() == Vector.class)) {
JSONArray array = (JSONArray) obj;
List<?> list = array.toJavaList((Class<?>) types[0]);
if (type == Set.class) {
Set<Object> set = Sets.newHashSet();
set.addAll(list);
return set;
} else if (type == Vector.class) {
Vector<Object> vec = new Vector<Object>(list.size());
vec.addAll(list);
return vec;
} else {
return list;
}
} else if (type == Map.class && obj instanceof JSONObject) {
JSONObject json = (JSONObject) obj;
Map<String, Object> map = Maps.newLinkedHashMap();
for (String key : json.keySet()) {
map.put(key, parse(types[1], json.get(key)));
}
return map;
} else {
log.warn("unsupport the type auto decode:{},{}", type, obj);
// TODO not support yet
return null;
}
} else {
Class<?> clz = (Class<?>) type;
if (type == String.class) {
return obj.toString();
} else if (type == Double.class) {
return Double.valueOf(obj.toString());
} else if (type == Integer.class) {
return Integer.valueOf(obj.toString());
} else if (type == Long.class) {
return Long.valueOf(obj.toString());
} else if (type == Boolean.class) {
return obj;
} else if (type == Short.class) {
return Short.valueOf(obj.toString());
} else if (type == Float.class) {
return Float.valueOf(obj.toString());
} else if (type == Character.class) {
return obj;
} else if (clz.isArray()) {
return (JSONArray) obj;
} else if (obj instanceof JSONObject) {
return JSON.toJavaObject((JSONObject) obj, clz);
} else if (type == Object.class){
return obj;
} else {
// TODO not support yet
log.debug("unsupport the type auto decode:{},{}", type, obj);
return obj;
}
}
} else {
return null;
}
}
private static Object parse(Class<?> type, Object obj) {
if (null != type && null != obj) {
if (type == String.class) {
return obj.toString();
} else if (type == Double.class) {
return Double.valueOf(obj.toString());
} else if (type == Integer.class) {
return Integer.valueOf(obj.toString());
} else if (type == Long.class) {
return Long.valueOf(obj.toString());
} else if (type == Boolean.class) {
return obj;
} else if (type == Short.class) {
return Short.valueOf(obj.toString());
} else if (type == Float.class) {
return Float.valueOf(obj.toString());
} else if (type == Character.class) {
return obj;
} else if (type.isArray()) {
return (JSONArray) obj;
} else if (type == Set.class) {
return (JSONArray) obj;
} else if (type == List.class) {
return (JSONArray) obj;
} else if (type == Map.class) {
return (JSONObject) obj;
} else if (type == Vector.class) {
return (JSONArray) obj;
} else {
if (type == Object.class) {
return obj;
} else {
return JSON.toJavaObject((JSONObject) obj, type);
}
}
}
return null;
}
public static Put buildPut(Object t) {
return buildPut(t, ReflectUtil.listField(t.getClass()));
}
public static Put buildPut(Object t, Field[] fields) {
String key = null;
for (Field field : ReflectUtil.listField(t.getClass())) {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
if (hf.isKey()) {
try {
key = field.get(t).toString();
} catch (Exception e) {
log.warn("Get key error:{}", t.getClass());
break;
}
}
}
}
if (!Strings.isNullOrEmpty(key)) {
try {
Put put = new Put(Bytes.toBytes(key));
for (Field field : ReflectUtil.listField(t.getClass())) {
HField hf = field.getAnnotation(HField.class);
Object v = field.get(t);
if (null != hf && null != v) {
String family = hf.family();
String qualifier = Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
byte[] f = Bytes.toBytes(family);
byte[] q = Bytes.toBytes(qualifier);
if (field.getType() == byte[].class) {
put.addColumn(f, q, (byte[]) v);
} else if (field.getType() == String.class) {
put.addColumn(f, q, Bytes.toBytes((String) v));
} else if (field.getType() == Double.class) {
put.addColumn(f, q, Bytes.toBytes((Double) v));
} else if (field.getType() == Integer.class) {
put.addColumn(f, q, Bytes.toBytes((Integer) v));
} else if (field.getType() == Long.class) {
put.addColumn(f, q, Bytes.toBytes((Long) v));
} else if (field.getType() == Date.class) {
if (null != field.getAnnotation(JsonFormat.class)
&& !Strings.isNullOrEmpty(field.getAnnotation(JsonFormat.class).pattern())) {
put.addColumn(f, q, Bytes.toBytes(
TimeUtil.format((Date) v, field.getAnnotation(JsonFormat.class).pattern())));
} else {
put.addColumn(f, q, Bytes.toBytes(TimeUtil.format((Date) v)));
}
} else if (field.getType() == Boolean.class) {
put.addColumn(f, q, Bytes.toBytes((Boolean) v));
} else if (field.getType() == Short.class) {
put.addColumn(f, q, Bytes.toBytes((Short) v));
} else if (field.getType() == Float.class) {
put.addColumn(f, q, Bytes.toBytes((Float) v));
} else if (field.getType() == Character.class) {
put.addColumn(f, q, Bytes.toBytes((Character) v));
} else if (field.getType().isArray()) {
put.addColumn(f, q, Bytes.toBytes(JSON.toJSONString(v)));
} else if (field.getType() == Set.class) {
put.addColumn(f, q, Bytes.toBytes(JSON.toJSONString(v)));
} else if (field.getType() == List.class) {
put.addColumn(f, q, Bytes.toBytes(JSON.toJSONString(v)));
} else if (field.getType() == Map.class) {
put.addColumn(f, q, Bytes.toBytes(JSON.toJSONString(v)));
} else if (field.getType() == Vector.class) {
put.addColumn(f, q, Bytes.toBytes(JSON.toJSONString(v)));
} else {
put.addColumn(f, q, Bytes.toBytes(JSON.toJSONString(v)));
}
}
}
return put;
} catch (Exception e) {
e.printStackTrace();
return null;
}
} else {
log.warn("Key can not be null");
return null;
}
}
public static void check(Class<?> cls) {
String key = null;
Map<String, Set<String>> fields = Maps.newHashMap();
for (Field field : ReflectUtil.listField(cls)) {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
if (hf.isKey()) {
key = field.getName();
} else {
String qualifier = Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
String f = hf.family() + ":" + qualifier;
if (!fields.containsKey(f)) {
fields.put(f, Sets.newHashSet());
}
fields.get(f).add(field.getName());
}
}
}
if (Strings.isNullOrEmpty(key)) {
throw new RuntimeException("key is not defined");
}
List<Set<String>> conflicts = fields.entrySet().stream().filter(x -> x.getValue().size() > 1)
.map(x -> x.getValue()).collect(Collectors.toList());
if (!conflicts.isEmpty()) {
throw new RuntimeException("field is conflict:" + conflicts);
}
}
@SafeVarargs
public static <T> void addGetColumn(Get get, AFunction<T, ?>... functions) {
if (Check.check(functions)) {
for (AFunction<T, ?> func : functions) {
addGetColumn(get, func);
}
}
}
@SafeVarargs
public static <T> void negGetColumn(Get get, Class<T> cls, AFunction<T, ?>... functions) {
if (null != get && null != cls && Check.check(functions)) {
Map<String, Set<String>> negs = getNegColumn(cls, functions);
for (Field field : ReflectUtil.listField(cls)) {
HField hf = field.getAnnotation(HField.class);
if (null != hf) {
if (!hf.isKey() && negs.containsKey(hf.family())) {
String qualifier = Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
if (!negs.get(hf.family()).contains(qualifier)) {
get.addColumn(Bytes.toBytes(hf.family()), Bytes.toBytes(qualifier));
}
}
}
}
}
}
private static <T> Map<String, Set<String>> getNegColumn(Class<T> cls, AFunction<T, ?>... functions) {
Map<String, Set<String>> negs = Maps.newHashMap();
for (AFunction<T, ?> func : functions) {
String family = getFamily(func);
String qualifier = getQualifier(func);
if (!Strings.isNullOrEmpty(family) && !Strings.isNullOrEmpty(qualifier)) {
if (!negs.containsKey(family)) {
negs.put(family, Sets.newHashSet());
}
negs.get(family).add(qualifier);
}
}
return negs;
}
private static <T> void addGetColumn(Get get, AFunction<T, ?> fn) {
Field field = getHField(fn);
if (null != field) {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
if (!hf.isKey()) {
String qualifier = Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
get.addColumn(Bytes.toBytes(hf.family()), Bytes.toBytes(qualifier));
}
}
}
}
@SafeVarargs
public static <T> void addScanColumn(Scan scan, AFunction<T, ?>... functions) {
if (Check.check(functions)) {
for (AFunction<T, ?> func : functions) {
addScanColumn(scan, func);
}
}
}
public static <T> String getFamily(AFunction<T, ?> fn) {
Field field = getHField(fn);
if (null != field) {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
if (!hf.isKey()) {
return hf.family();
}
}
}
return null;
}
public static <T> String getQualifier(AFunction<T, ?> fn) {
Field field = getHField(fn);
if (null != field) {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
if (!hf.isKey()) {
return Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
}
}
}
return null;
}
private static <T> void addScanColumn(Scan scan, AFunction<T, ?> fn) {
Field field = getHField(fn);
if (null != field) {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
if (!hf.isKey()) {
String qualifier = Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
scan.addColumn(Bytes.toBytes(hf.family()), Bytes.toBytes(qualifier));
}
}
}
}
public static <T> void negScanColumn(Scan scan, Class<T> cls, AFunction<T, ?>... functions) {
if (null != scan && null != cls && Check.check(functions)) {
Map<String, Set<String>> negs = getNegColumn(cls, functions);
for (Field field : ReflectUtil.listField(cls)) {
HField hf = field.getAnnotation(HField.class);
if (null != hf) {
if (!hf.isKey() && negs.containsKey(hf.family())) {
String qualifier = Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
if (!negs.get(hf.family()).contains(qualifier)) {
scan.addColumn(Bytes.toBytes(hf.family()), Bytes.toBytes(qualifier));
}
}
}
}
}
}
private static <T> Field getHField(AFunction<T, ?> fn) {
Field field = FIELD_CACHE.get(fn.getClass());
if (null == field) {
Method writeReplaceMethod;
try {
writeReplaceMethod = fn.getClass().getDeclaredMethod("writeReplace");
} catch (NoSuchMethodException e) {
throw new RuntimeException(e);
}
boolean isAccessible = writeReplaceMethod.isAccessible();
writeReplaceMethod.setAccessible(true);
SerializedLambda serializedLambda;
try {
serializedLambda = (SerializedLambda) writeReplaceMethod.invoke(fn);
} catch (IllegalAccessException | InvocationTargetException e) {
throw new RuntimeException(e);
}
writeReplaceMethod.setAccessible(isAccessible);
String fieldName = serializedLambda.getImplMethodName().substring("get".length());
fieldName = fieldName.replaceFirst(fieldName.charAt(0) + "", (fieldName.charAt(0) + "").toLowerCase());
try {
field = Class.forName(serializedLambda.getImplClass().replace("/", ".")).getDeclaredField(fieldName);
} catch (ClassNotFoundException | NoSuchFieldException e) {
throw new RuntimeException(e);
}
if (null != field) {
FIELD_CACHE.put(fn.getClass(), field);
}
}
return field;
}
public static String generateSchema(Class<?> cls) {
StringBuilder sb = new StringBuilder();
String key = null;
Map<String, Set<String>> fields = Maps.newHashMap();
for (Field field : ReflectUtil.listField(cls)) {
HField hf = field.getAnnotation(HField.class);
if (hf != null) {
sb.append("|" + field.getName() + "|" + field.getType().getName() + "|");
if (hf.isKey()) {
sb.append("-|-|true");
} else {
String qualifier = Strings.isNullOrEmpty(hf.qualifier()) ? field.getName() : hf.qualifier();
String f = hf.family();
sb.append(f + "|" + qualifier + "|-|");
}
sb.append("\n");
}
}
return sb.toString();
}
}
Hbase操作类
接口
public interface HbaseOperations {
public <T> List<T> scan(Scan scan, Class<T> cls);
public <T> List<T> scan(String tableName, Scan scan, final RowMapper<T> rowmapper);
/**
* 遍历所有符合条件的数据,慎用
* @param <T>
* @param tableName
* @param scan
* @param rowmapper
*/
public <T> void find(String tableName, Scan scan, final RowMapper<T> rowmapper);
public <T> List<T> scan(String tableName, Scan scan, final RowMapper<T> rowmapper, long maxRows);
public <T> T get(Get get, Class<T> cls);
public <T> List<T> get(List<Get> gets, Class<T> cls);
public <T> List<T> get(String tableName, List<Get> gets, Class<T> cls);
public <T> T get(String tableName, Get get, final RowMapper<T> mapper);
public <T> List<T> get(String tableName, List<Get> gets, final RowMapper<T> mapper);
public <T> T get(String goodsHbaseTable, String key, String family, final RowMapper<T> rowMapper);
public <T> T get(String tableName, String key, String family, String qualifier, RowMapper<T> mapper);
public <T> void put(T obj);
public void put(String tableName, Put put);
public <T> void put(Put put, Class<T> cls);
public void putAll(String tableName, List<Put> puts);
public <T> void putAll(List<T> objs);
public void putAll(Class<?> cls, List<Put> puts);
public <T> void deleteObject(T obj);
public <T> void deleteObject(String key, Class<T> cls);
public void deleteRow(String tableName, String key);
public void delete(String tableName, Delete delete);
public void delete(String tableName, List<Delete> deletes);
public void deleteRows(String tableName, List<String> keys);
public <T> void deleteObjects(List<T> objs);
}
操作实现类
public class HbaseTemplate implements HbaseOperations {
private static final Logger LOGGER = LoggerFactory.getLogger(HbaseTemplate.class);
private static final int DEFAULT_SCAN_NUM = 100;
private Configuration configuration;
private volatile Connection connection;
public HbaseTemplate(Configuration configuration) {
this.setConfiguration(configuration);
Assert.notNull(configuration, " a valid configuration is required");
}
private <T> T execute(String tableName, TableCallback<T> action) {
Assert.notNull(action, "Callback object must not be null");
Assert.notNull(tableName, "No table specified");
StopWatch sw = new StopWatch();
sw.start();
Table table = null;
try {
table = this.getConnection().getTable(TableName.valueOf(tableName));
return action.doInTable(table);
} catch (Throwable throwable) {
throw new HbaseSystemException(throwable);
} finally {
if (null != table) {
try {
table.close();
sw.stop();
} catch (IOException e) {
LOGGER.error("hbase资源释放失败");
}
}
}
}
@Override
public <T> List<T> get(String tableName, List<Get> gets, final RowMapper<T> mapper) {
return execute(tableName, new TableCallback<List<T>>() {
@Override
public List<T> doInTable(Table table) throws Throwable {
List<T> list = Lists.newArrayList();
Result[] rs = table.get(gets);
for (int i = 0; i < rs.length; i++) {
list.add(mapper.mapRow(rs[i], i));
}
return list;
}
});
}
@Override
public <T> List<T> get(List<Get> gets, Class<T> cls) {
Htable htable = cls.getAnnotation(Htable.class);
if (null != htable) {
return get(htable.table(), gets, cls);
}
return Lists.newArrayList();
}
@Override
public <T> List<T> get(String table, List<Get> gets, Class<T> cls) {
if (!Strings.isNullOrEmpty(table)) {
return get(table, gets, new RowMapper<T>() {
@Override
public T mapRow(Result result, int rowNum) throws Exception {
T t = cls.newInstance();
HMapper.mapper(result, t);
return t;
}
});
}
return Lists.newArrayList();
}
@Override
public <T> T get(String tableName, Get get, final RowMapper<T> mapper) {
return execute(tableName, new TableCallback<T>() {
@Override
public T doInTable(Table table) throws Throwable {
Result r = table.get(get);
if (!r.isEmpty()) {
return mapper.mapRow(r, 0);
}
return null;
}
});
}
@Override
public <T> T get(Get get, Class<T> cls) {
Htable htable = cls.getAnnotation(Htable.class);
if (null != htable && !Strings.isNullOrEmpty(htable.table())) {
return get(htable.table(), get, new RowMapper<T>() {
@Override
public T mapRow(Result result, int rowNum) throws Exception {
T t = cls.newInstance();
HMapper.mapper(result, t);
return t;
}
});
}
return null;
}
@Override
public <T> T get(String tableName, String key, String family, RowMapper<T> mapper) {
return execute(tableName, new TableCallback<T>() {
@Override
public T doInTable(Table table) throws Throwable {
Get get = new Get(Bytes.toBytes(key));
get.addFamily(Bytes.toBytes(family));
Result r = table.get(get);
if (!r.isEmpty()) {
return mapper.mapRow(r, 0);
}
return null;
}
});
}
@Override
public <T> T get(String tableName, String key, String family, String qualifier, RowMapper<T> mapper) {
return execute(tableName, new TableCallback<T>() {
@Override
public T doInTable(Table table) throws Throwable {
Get get = new Get(Bytes.toBytes(key));
get.addColumn(Bytes.toBytes(family), Bytes.toBytes(qualifier));
Result r = table.get(get);
if (!r.isEmpty()) {
return mapper.mapRow(r, 0);
}
return null;
}
});
}
@Override
public <T> List<T> scan(Scan scan, Class<T> cls) {
if (null != scan && null != cls) {
Htable htable = cls.getAnnotation(Htable.class);
if (null != htable && !Strings.isNullOrEmpty(htable.table())) {
return scan(htable.table(), scan, new RowMapper<T>() {
@Override
public T mapRow(Result result, int rowNum) throws Exception {
T t = cls.newInstance();
HMapper.mapper(result, t);
return t;
}
}, DEFAULT_SCAN_NUM);
}
}
return Lists.newArrayList();
}
@Override
public <T> List<T> scan(String tableName, Scan scan, RowMapper<T> rowMapper) {
return scan(tableName, scan, rowMapper, DEFAULT_SCAN_NUM);
}
@Override
public <T> List<T> scan(String tableName, Scan scan, RowMapper<T> rowMapper, long maxRows) {
if (maxRows > 0) {
if(scan.getMaxResultSize() == -1) {
scan.setMaxResultSize(maxRows);
}
if(scan.getCaching() == -1) {
scan.setCaching(1000);
}
return execute(tableName, new TableCallback<List<T>>() {
@Override
public List<T> doInTable(Table table) throws Throwable {
List<T> result = Lists.newLinkedList();
ResultScanner rs = table.getScanner(scan);
int i = 0;
Iterator<Result> it = rs.iterator();
while (it.hasNext()) {
Result r = it.next();
T t = rowMapper.mapRow(r, i++);
if (null != t) {
result.add(t);
}
if (i >= scan.getMaxResultSize()) {
break;
}
}
return result;
}
});
} else {
return Lists.newArrayList();
}
}
@Override
public void deleteRow(String tableName, String key) {
execute(tableName, new TableCallback<Object>() {
@Override
public Object doInTable(Table htable) throws Throwable {
htable.delete(new Delete(Bytes.toBytes(key)));
return null;
}
});
}
@Override
public void deleteRows(String tableName, List<String> keys) {
execute(tableName, new TableCallback<Object>() {
@Override
public Object doInTable(Table htable) throws Throwable {
List<Delete> dels = Lists.newArrayList();
for (String key : keys) {
Delete del = new Delete(Bytes.toBytes(key));
dels.add(del);
}
htable.delete(dels);
return null;
}
});
}
@Override
public <T> void deleteObjects(List<T> objs) {
if (null != objs && !objs.isEmpty()) {
Class<?> cls = objs.get(0).getClass();
Htable htable = cls.getAnnotation(Htable.class);
Field kf = getKeyField(cls.getClass());
if (null != htable && !Strings.isNullOrEmpty(htable.table()) && null != kf) {
String tableName = htable.table();
List<Delete> deletes = Lists.newArrayList();
for (T t : objs) {
try {
Delete del = new Delete(Bytes.toBytes(kf.get(t).toString()));
deletes.add(del);
} catch (Exception e) {
e.printStackTrace();
}
}
if (Check.check(deletes)) {
delete(tableName, deletes);
}
}
}
}
@Override
public <T> void deleteObject(T obj) {
Htable htable = obj.getClass().getAnnotation(Htable.class);
if (htable != null) {
try {
Field kf = getKeyField(obj.getClass());
Object key = kf.get(obj);
if (key != null) {
deleteRow(htable.table(), key.toString());
}
} catch (Exception e) {
e.printStackTrace();
}
}
}
@Override
public <T> void deleteObject(String key, Class<T> cls) {
Htable htable = cls.getAnnotation(Htable.class);
if (htable != null) {
deleteRow(htable.table(), key);
}
}
@Override
public void delete(String tableName, Delete delete) {
if (!Strings.isNullOrEmpty(tableName)) {
execute(tableName, new TableCallback<Object>() {
@Override
public Object doInTable(Table table) throws Throwable {
table.delete(delete);
return null;
}
});
}
}
@Override
public void delete(String tableName, List<Delete> deletes) {
if (!Strings.isNullOrEmpty(tableName) && Check.check(deletes)) {
execute(tableName, new TableCallback<Object>() {
@Override
public Object doInTable(Table table) throws Throwable {
table.delete(deletes);
return null;
}
});
}
}
private <T> Field getKeyField(Class<T> cls) {
for (Field field : ReflectUtil.listField(cls)) {
HField hf = field.getAnnotation(HField.class);
if (null != hf && hf.isKey()) {
return field;
}
}
return null;
}
@Override
public void putAll(String tableName, List<Put> puts) {
execute(tableName, new TableCallback<Object>() {
@Override
public Object doInTable(Table htable) throws Throwable {
htable.put(puts);
return null;
}
});
}
@Override
public void putAll(Class<?> cls, List<Put> puts) {
Htable htable = cls.getAnnotation(Htable.class);
if (null != htable) {
putAll(htable.table(), puts);
}
}
@Override
public void put(String tableName, Put put) {
execute(tableName, new TableCallback<Object>() {
@Override
public Object doInTable(Table htable) throws Throwable {
htable.put(put);
return null;
}
});
}
@Override
public <T> void put(Put put, Class<T> cls) {
Htable htable = cls.getAnnotation(Htable.class);
if (htable != null) {
put(htable.table(), put);
}
}
@Override
public <T> void put(T obj) {
if (null != obj) {
if (obj instanceof List) {
putAll((List<?>) obj);
} else {
Class<?> cls = obj.getClass();
Htable htable = cls.getAnnotation(Htable.class);
if (null != htable && !Strings.isNullOrEmpty(htable.table())) {
try {
Put put = HMapper.buildPut(obj);
put(htable.table(), put);
} catch (Exception e) {
e.printStackTrace();
}
}
}
}
}
public <T> void putAll(List<T> objs) {
if (null != objs && !objs.isEmpty()) {
Class<?> cls = objs.get(0).getClass();
Htable htable = cls.getAnnotation(Htable.class);
if (null != htable && !Strings.isNullOrEmpty(htable.table())) {
List<Put> puts = Lists.newArrayList();
for (T t : objs) {
try {
Put put = HMapper.buildPut(t);
if (null != put) {
puts.add(put);
}
} catch (Exception e) {
e.printStackTrace();
}
}
if (!puts.isEmpty()) {
putAll(htable.table(), puts);
}
}
}
}
public void setConnection(Connection connection) {
this.connection = connection;
}
public Connection getConnection() {
if (null == this.connection) {
synchronized (this) {
if (null == this.connection) {
try {
ThreadPoolExecutor poolExecutor = new ThreadPoolExecutor(200, Integer.MAX_VALUE, 60L,
TimeUnit.SECONDS, new SynchronousQueue<Runnable>());
// init pool
poolExecutor.prestartCoreThread();
this.connection = ConnectionFactory.createConnection(configuration, poolExecutor);
} catch (IOException e) {
LOGGER.error("site.assad.spring.data.hbase connection资源池创建失败");
}
}
}
}
return this.connection;
}
public Configuration getConfiguration() {
return configuration;
}
public void setConfiguration(Configuration configuration) {
this.configuration = configuration;
}
@Override
public <T> void find(String tableName, Scan scan, RowMapper<T> rowMapper) {
if(scan.getCaching() == -1) {
scan.setCaching(1000);
}
execute(tableName, new TableCallback<List<T>>() {
@Override
public List<T> doInTable(Table table) throws Throwable {
ResultScanner rs = table.getScanner(scan);
int i = 0;
Iterator<Result> it = rs.iterator();
while (it.hasNext()) {
Result r = it.next();
rowMapper.mapRow(r, i++);
}
return Lists.newArrayList();
}
});
}
}
获取Hbase操作类
获取类
@Configuration
@ConditionalOnClass(HbaseTemplate.class)
public class HbaseAutoConfiguration {
@Autowired(required = false)
private HbaseProperties hbaseProperties;
@Bean
@Primary
@ConditionalOnMissingBean(HbaseTemplate.class)
public HbaseTemplate hbaseTemplate() {
return new HbaseTemplate(getHbaseConfiguration());
}
private org.apache.hadoop.conf.Configuration getHbaseConfiguration() {
org.apache.hadoop.conf.Configuration configuration = HBaseConfiguration.create();
if (hbaseProperties == null) {
return configuration;
}
hbaseProperties.buildConfiguration(configuration);
return configuration;
}
}
反射工具类
public final class ReflectUtil {
private ReflectUtil() {
}
private static HashMap<Class<?>, Field[]> fieldsHolder = new HashMap<Class<?>, Field[]>();
private static HashMap<Class<?>, Map<Class<?>, Field[]>> anaFieldsHolder = new HashMap<Class<?>, Map<Class<?>, Field[]>>();
private static HashMap<Class<?>, Map<Class<?>, Field[]>> unAnaFieldsHolder = new HashMap<Class<?>, Map<Class<?>, Field[]>>();
private static HashMap<Class<?>, Map<Class<?>, Map<Field, Annotation>>> anaCacheHolder = new HashMap<Class<?>, Map<Class<?>, Map<Field, Annotation>>>();
public static Field[] listField(Class<?> c) {
if (!fieldsHolder.containsKey(c)) {
initField(c);
}
return fieldsHolder.get(c);
}
private synchronized static void initField(Class<?> c) {
if (!fieldsHolder.containsKey(c)) {
List<Field> fieldList = new ArrayList<Field>();
Class<?> ts = c;
while (ts != null) {
for (Field field : ts.getDeclaredFields()) {
if (!Modifier.isStatic(field.getModifiers())) {
field.setAccessible(true);
fieldList.add(field);
}
}
ts = ts.getSuperclass();
}
fieldsHolder.put(c, fieldList.stream().toArray(Field[]::new));
}
}
public static <T extends Annotation> Field[] listField(Class<?> c, Class<T> annotationCls) {
if (!anaFieldsHolder.containsKey(c)) {
initField(c, annotationCls);
}
return anaFieldsHolder.get(c).get(annotationCls);
}
private synchronized static <T extends Annotation> void initField(Class<?> c, Class<T> annotationCls) {
if (!anaFieldsHolder.containsKey(c)) {
anaFieldsHolder.put(c, new HashMap<Class<?>, Field[]>());
List<Field> fieldList = new ArrayList<Field>();
Class<?> ts = c;
while (ts != null) {
for (Field field : ts.getDeclaredFields()) {
if (!Modifier.isStatic(field.getModifiers()) && null != field.getAnnotation(annotationCls)) {
field.setAccessible(true);
fieldList.add(field);
}
}
ts = ts.getSuperclass();
}
anaFieldsHolder.get(c).put(annotationCls, fieldList.stream().toArray(Field[]::new));
}
}
public static <T extends Annotation> Field[] listUnAnnoField(Class<?> c, Class<T> annotationCls) {
if (!unAnaFieldsHolder.containsKey(c)) {
initUnAnnoField(c, annotationCls);
}
return unAnaFieldsHolder.get(c).get(annotationCls);
}
private synchronized static <T extends Annotation> void initUnAnnoField(Class<?> c, Class<T> annotationCls) {
if (!unAnaFieldsHolder.containsKey(c)) {
unAnaFieldsHolder.put(c, new HashMap<Class<?>, Field[]>());
List<Field> fieldList = new ArrayList<Field>();
Class<?> ts = c;
while (ts != null) {
for (Field field : ts.getDeclaredFields()) {
if (!Modifier.isFinal(field.getModifiers()) && !Modifier.isStatic(field.getModifiers()) && null == field.getAnnotation(annotationCls)) {
field.setAccessible(true);
fieldList.add(field);
}
}
ts = ts.getSuperclass();
}
unAnaFieldsHolder.get(c).put(annotationCls, fieldList.stream().toArray(Field[]::new));
}
}
@SuppressWarnings("unchecked")
public static <T extends Annotation> T getAnnotation(Class<?>c, Field field, Class<T> annotationCls) {
if(!anaCacheHolder.containsKey(annotationCls) || !anaCacheHolder.get(annotationCls).containsKey(c)) {
initAnnotationOfClass(c, annotationCls);
}
return (T) anaCacheHolder.get(annotationCls).get(c).get(field);
}
private synchronized static <T extends Annotation> void initAnnotationOfClass(Class<?> c, Class<T> annotationCls) {
if(!anaCacheHolder.containsKey(annotationCls)) {
anaCacheHolder.put(annotationCls, new HashMap<Class<?>, Map<Field, Annotation>>());
}
Field[] fields = listField(c, annotationCls);
Map<Field, Annotation> annos = new HashMap<Field, Annotation>();
for(Field field : fields) {
T anno = field.getAnnotation(annotationCls);
annos.put(field, anno);
}
anaCacheHolder.get(annotationCls).put(c, annos);
}
}
bean 配置
<bean name="hbaseAutoConfiguration" class="HbaseAutoConfiguration">
</bean>