【生意多】-免费发布分类信息
当前位置: 首页 » 新闻 » 教程 » 操作系统 » 正文

鏁村悎鍘熺悊妗堜緥(鏁村悎璁捐妗堜緥)

放大字体  缩小字体 发布日期:2022-06-30 22:40:33    浏览次数:21

*要求:读取Hbase当中user这张表的f1:name、f1:age数据,将数据写入到另外一张user2表的f1列族里面去==****

第一步:创建表

注意:两张表的列族一定要相同

第二步:创建maven工程并导入jar包

pom.xml文件内容如下:

<?xml version="1.0" encoding="UTF-8"?><project xmlns="http://maven.apache.org/POM/4.0.0"         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">    <modelVersion>4.0.0</modelVersion>    <groupId>Hadoop</groupId>    <artifactId>HbaseTang</artifactId>    <version>1.0-SNAPSHOT</version>    <packaging>jar</packaging>    <properties>        <hadoop.version>2.7.3</hadoop.version>    </properties>    <dependencies>        <dependency>            <groupId>commons-cli</groupId>            <artifactId>commons-cli</artifactId>            <version>1.2</version>        </dependency>        <dependency>            <groupId>commons-logging</groupId>            <artifactId>commons-logging</artifactId>            <version>1.1.3</version>        </dependency>        <dependency>            <groupId>org.apache.hadoop</groupId>            <artifactId>hadoop-mapreduce-client-jobclient</artifactId>            <version>${hadoop.version}</version>        </dependency>        <dependency>            <groupId>org.apache.hadoop</groupId>            <artifactId>hadoop-common</artifactId>            <version>${hadoop.version}</version>        </dependency>        <dependency>            <groupId>org.apache.hadoop</groupId>            <artifactId>hadoop-hdfs</artifactId>            <version>2.7.3</version>        </dependency>        <dependency>            <groupId>org.apache.hadoop</groupId>            <artifactId>hadoop-hdfs</artifactId>            <version>${hadoop.version}</version>        </dependency>        <dependency>            <groupId>org.apache.hadoop</groupId>            <artifactId>hadoop-mapreduce-client-app</artifactId>            <version>${hadoop.version}</version>        </dependency>        <dependency>            <groupId>org.apache.hadoop</groupId>            <artifactId>hadoop-mapreduce-client-hs</artifactId>            <version>${hadoop.version}</version>        </dependency>        <dependency>            <groupId>org.apache.hbase</groupId>            <artifactId>hbase-client</artifactId>            <version>1.2.1</version>        </dependency>        <dependency>            <groupId>org.apache.hbase</groupId>            <artifactId>hbase-common</artifactId>            <version>1.2.1</version>        </dependency>        <dependency>            <groupId>org.apache.hbase</groupId>            <artifactId>hbase-server</artifactId>            <version>1.2.1</version>        </dependency>        <dependency>            <groupId>junit</groupId>            <artifactId>junit</artifactId>            <version>4.12</version>        </dependency>    </dependencies></project>

第三步:开发MR程序实现功能

(1)自定义map类

import org.apache.hadoop.hbase.Cell;
import org.apache.hadoop.hbase.CellUtil;
import org.apache.hadoop.hbase.client.Put;
import org.apache.hadoop.hbase.client.Result;
import org.apache.hadoop.hbase.io.ImmutableBytesWritable;
import org.apache.hadoop.hbase.mapreduce.TableMapper;
import org.apache.hadoop.hbase.util.Bytes;
import org.apache.hadoop.io.Text;

import java.io.IOException;


public class HbaseReadMapper extends TableMapper&lt;Text, Put&gt; {
`
`@Override`
protected void map(ImmutableBytesWritable key, Result value, Context context) throws IOException, InterruptedException {

  • 			
     
  • //获得roweky的字节数组
    byte[] rowkey_bytes = key.get();
    String rowkeyStr = Bytes.toString(rowkey_bytes);
    Text text = new Text(rowkeyStr);

    //输出数据 -&gt; 写数据 -&gt; Put 构建Put对象
    Put put = new Put(rowkey_bytes);
    //获取一行中所有的Cell对象
    Cell[] cells = value.rawCells();
    //将f1 : name& age输出
    for(Cell cell: cells) {
    //当前cell是否是f1
    //列族
    byte[] family_bytes = CellUtil.cloneFamily(cell);
    String familyStr = Bytes.toString(family_bytes);
    if("f1".equals(familyStr)) {
    //在判断是否是name | age
    byte[] qualifier_bytes = CellUtil.cloneQualifier(cell);
    String qualifierStr = Bytes.toString(qualifier_bytes);
    if("name".equals(qualifierStr)) {
    put.add(cell);
    }
    if("age".equals(qualifierStr)) {
    put.add(cell);
    }
    }
    }

    //判断是否为空;不为空,才输出
    if(!put.isEmpty()){
    context.write(text, put);
    }
    }
    }

  • (2)自定义reduce类

    package com.kaikeba.hbase.demo01;import org.apache.hadoop.hbase.client.Put;import org.apache.hadoop.hbase.io.ImmutableBytesWritable;import org.apache.hadoop.hbase.mapreduce.TableReducer;import org.apache.hadoop.io.Text;import java.io.IOException;public class HbaseWriteReducer extends TableReducer<Text, Put, ImmutableBytesWritable> {    //将map传输过来的数据,写入到hbase表        @Override    protected void reduce(Text key, Iterable<Put> values, Context context) throws IOException, InterruptedException {                //rowkey        ImmutableBytesWritable immutableBytesWritable = new ImmutableBytesWritable();        immutableBytesWritable.set(key.toString().getBytes());        //遍历put对象,并输出        for(Put put: values) {            context.write(immutableBytesWritable, put);        }    }}
    •  

    (3)main入口类

    import org.apache.hadoop.conf.Configuration;import org.apache.hadoop.conf.Configured;import org.apache.hadoop.hbase.HbaseConfiguration;import org.apache.hadoop.hbase.TableName;import org.apache.hadoop.hbase.client.Put;import org.apache.hadoop.hbase.client.Scan;import org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil;import org.apache.hadoop.io.Text;import org.apache.hadoop.mapreduce.Job;import org.apache.hadoop.util.Tool;import org.apache.hadoop.util.ToolRunner;public class HbaseMR extends Configured implements Tool {    public static void main(String[] args) throws Exception {        Configuration configuration = HbaseConfiguration.create();        //设定绑定的zk集群        configuration.set("hbase.zookeeper.quorum", "node01:2181,node02:2181,node03:2181");        int run = ToolRunner.run(configuration, new HbaseMR(), args);        System.exit(run);    }    @Override    public int run(String[] args) throws Exception {        Job job = Job.getInstance(super.getConf());        job.setJarByClass(HbaseMR.class);        //mapper        TableMapReduceUtil.initTableMapperJob(TableName.valueOf("myuser"), new Scan(),HbaseReadMapper.class, Text.class, Put.class, job);        //reducer        TableMapReduceUtil.initTableReducerJob("myuser2", HbaseWriteReducer.class, job);        boolean b = job.waitForCompletion(true);        return b? 0: 1;    }}

    第四步:打成jar包提交到集群运行

    打包:


    执行命令:

    hadoop jar HbaseTang-1.0-SNAPSHOT.jar mapreduce_hbase.HbaseMR

    执行结果:

     
    (文/小编)
    打赏
    免责声明
    • 
    本文为小编原创作品,作者: 小编。欢迎转载,转载请注明原文出处:https://www.31duo.com/news/show-3574857.html 。本文仅代表作者个人观点,本站未对其内容进行核实,请读者仅做参考,如若文中涉及有违公德、触犯法律的内容,一经发现,立即删除,作者需自行承担相应责任。涉及到版权或其他问题,请及时联系我们。
     

    (c)2016-2019 31DUO.COM All Rights Reserved浙ICP备19001410号-4

    浙ICP备19001410号-4