MapReduce-从HBase读取处理后再写入HBase
代码如下
1package com.hbase.mapreduce; 2 3import java.io.IOException; 4 5import org.apache.hadoop.conf.Configuration; 6import org.apache.hadoop.conf.Configured; 7import org.apache.hadoop.hbase.Cell; 8import org.apache.hadoop.hbase.CellUtil; 9import org.apache.hadoop.hbase.HBaseConfiguration; 10import org.apache.hadoop.hbase.KeyValue; 11import org.apache.hadoop.hbase.client.Mutation; 12import org.apache.hadoop.hbase.client.Put; 13import org.apache.hadoop.hbase.client.Result; 14import org.apache.hadoop.hbase.client.Scan; 15import org.apache.hadoop.hbase.io.ImmutableBytesWritable; 16import org.apache.hadoop.hbase.mapreduce.TableInputFormat; 17import org.apache.hadoop.hbase.mapreduce.TableMapper; 18import org.apache.hadoop.hbase.mapreduce.TableOutputFormat; 19import org.apache.hadoop.hbase.mapreduce.TableReducer; 20import org.apache.hadoop.hbase.util.Bytes; 21import org.apache.hadoop.mapreduce.Job; 22import org.apache.hadoop.mapreduce.Mapper; 23import org.apache.hadoop.mapreduce.Reducer; 24import org.apache.hadoop.util.Tool; 25import org.apache.hadoop.util.ToolRunner; 26 27/** 28* @author:FengZhen 29* @create:2018年9月17日 30* 从HBase读写入HBase 31* zip -d HBaseToHBase.jar 'META-INF/.SF' 'META-INF/.RSA' 'META-INF/*SF' 32*/ 33public class HBaseToHBase extends Configured implements Tool{ 34 35 private static String addr="HDP233,HDP232,HDP231"; 36 private static String port="2181"; 37 38 public enum Counters { ROWS, COLS, VALID, ERROR, EMPTY, NOT_EMPTY} 39 40 static class ParseMapper extends TableMapper<ImmutableBytesWritable, Put>{ 41 private byte[] columnFamily = null; 42 @Override 43 protected void setup(Mapper<ImmutableBytesWritable, Result, ImmutableBytesWritable, Put>.Context context) 44 throws IOException, InterruptedException { 45 columnFamily = Bytes.toBytes(context.getConfiguration().get("conf.columnfamily")); 46 } 47 @Override 48 protected void map(ImmutableBytesWritable key, Result value, 49 Mapper<ImmutableBytesWritable, Result, ImmutableBytesWritable, Put>.Context context) 50 throws IOException, InterruptedException { 51 context.getCounter(Counters.ROWS).increment(1); 52 String hbaseValue = null; 53 54 Put put = new Put(key.get()); 55 for (Cell cell : value.listCells()) { 56 context.getCounter(Counters.COLS).increment(1); 57 hbaseValue = Bytes.toString(CellUtil.cloneValue(cell)); 58 if (hbaseValue.length() > 0) { 59 String top = hbaseValue.substring(0, hbaseValue.length()/2); 60 String detail = hbaseValue.substring(hbaseValue.length()/2, hbaseValue.length() - 1); 61 put.addColumn(columnFamily, Bytes.toBytes("top"), Bytes.toBytes(top)); 62 put.addColumn(columnFamily, Bytes.toBytes("detail"), Bytes.toBytes(detail)); 63 context.getCounter(Counters.NOT_EMPTY).increment(1); 64 }else { 65 put.addColumn(columnFamily, Bytes.toBytes("empty"), Bytes.toBytes(hbaseValue)); 66 context.getCounter(Counters.EMPTY).increment(1); 67 } 68 } 69 try { 70 context.write(key, put); 71 context.getCounter(Counters.VALID).increment(1); 72 } catch (Exception e) { 73 e.printStackTrace(); 74 context.getCounter(Counters.ERROR).increment(1); 75 } 76 } 77 } 78 79 static class ParseTableReducer extends TableReducer<ImmutableBytesWritable, Put, ImmutableBytesWritable>{ 80 @Override 81 protected void reduce(ImmutableBytesWritable key, Iterable<Put> values, 82 Reducer<ImmutableBytesWritable, Put, ImmutableBytesWritable, Mutation>.Context context) 83 throws IOException, InterruptedException { 84 for (Put put : values) { 85 context.write(key, put); 86 } 87 } 88 } 89 90 public int run(String[] arg0) throws Exception { 91 String table = arg0[0]; 92 String column = arg0[1]; 93 String destTable = arg0[2]; 94 95 Configuration configuration = HBaseConfiguration.create(); 96 configuration.set("hbase.zookeeper.quorum",addr); 97 configuration.set("hbase.zookeeper.property.clientPort", port); 98 99 Scan scan = new Scan(); 100 if (null != column) { 101 byte[][] colkey = KeyValue.parseColumn(Bytes.toBytes(column)); 102 if (colkey.length > 1) { 103 scan.addColumn(colkey[0], colkey[1]); 104 configuration.set("conf.columnfamily", Bytes.toString(colkey[0])); 105 configuration.set("conf.columnqualifier", Bytes.toString(colkey[1])); 106 }else { 107 scan.addFamily(colkey[0]); 108 configuration.set("conf.columnfamily", Bytes.toString(colkey[0])); 109 } 110 } 111 112 Job job = Job.getInstance(configuration); 113 job.setJobName("HBaseToHBase2"); 114 job.setJarByClass(HBaseToHBase2.class); 115 116 job.getConfiguration().set(TableInputFormat.INPUT_TABLE, table); 117 job.getConfiguration().set(TableOutputFormat.OUTPUT_TABLE, destTable); 118 119 job.setMapperClass(ParseMapper.class); 120 job.setMapOutputKeyClass(ImmutableBytesWritable.class); 121 job.setMapOutputValueClass(Put.class); 122 123// job.setReducerClass(ParseTableReducer.class); 124 job.setOutputKeyClass(ImmutableBytesWritable.class); 125 job.setOutputValueClass(Put.class); 126 127 job.setInputFormatClass(TableInputFormat.class); 128 TableInputFormat.addColumns(scan, KeyValue.parseColumn(Bytes.toBytes(column))); 129 job.setOutputFormatClass(TableOutputFormat.class); 130 131 job.setNumReduceTasks(0); 132 133 //使用TableMapReduceUtil会报类找不到错误 134 //Caused by: java.lang.ClassNotFoundException: com.yammer.metrics.core.MetricsRegistry 135// TableMapReduceUtil.initTableMapperJob(table, scan, ParseMapper.class, ImmutableBytesWritable.class, Put.class, job); 136// TableMapReduceUtil.initTableReducerJob(table, IdentityTableReducer.class, job); 137 138 return job.waitForCompletion(true) ? 0 : 1; 139 } 140 public static void main(String[] args) throws Exception { 141 String[] params = new String[] {"test_table_mr", "data:info", "test_table_dest"}; 142 int exitCode = ToolRunner.run(new HBaseToHBase2(), params); 143 System.exit(exitCode); 144 } 145}
打包测试
1zip -d HBaseToHBase.jar 'META-INF/.SF' 'META-INF/.RSA' 'META-INF/*SF' 2hadoop jar HBaseToHBase.jar com.hbase.mapreduce.HBaseToHBase
出现的问题
一开始使用额TableMapReduceUtil,但是报下面这个错
1Exception in thread "main" java.lang.NoClassDefFoundError: com/yammer/metrics/core/MetricsRegistry 2 at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.addHBaseDependencyJars(TableMapReduceUtil.java:732) 3 at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.addDependencyJars(TableMapReduceUtil.java:777) 4 at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.initTableMapperJob(TableMapReduceUtil.java:212) 5 at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.initTableMapperJob(TableMapReduceUtil.java:168) 6 at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.initTableMapperJob(TableMapReduceUtil.java:291) 7 at org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil.initTableMapperJob(TableMapReduceUtil.java:92) 8 at com.hbase.mapreduce.HBaseToHBase.run(HBaseToHBase.java:108) 9 at org.apache.hadoop.util.ToolRunner.run(ToolRunner.java:76) 10 at org.apache.hadoop.util.ToolRunner.run(ToolRunner.java:90) 11 at com.hbase.mapreduce.HBaseToHBase.main(HBaseToHBase.java:115) 12 at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method) 13 at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) 14 at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) 15 at java.lang.reflect.Method.invoke(Method.java:498) 16 at org.apache.hadoop.util.RunJar.run(RunJar.java:233) 17 at org.apache.hadoop.util.RunJar.main(RunJar.java:148) 18Caused by: java.lang.ClassNotFoundException: com.yammer.metrics.core.MetricsRegistry 19 at java.net.URLClassLoader.findClass(URLClassLoader.java:381) 20 at java.lang.ClassLoader.loadClass(ClassLoader.java:424) 21 at sun.misc.Launcher$AppClassLoader.loadClass(Launcher.java:338) 22 at java.lang.ClassLoader.loadClass(ClassLoader.java:357) 23 ... 16 more
解决,不使用TableMapReduceUtil,分布设置便可解决此问题