使用BulkLoad从HDFS批量导入数据到HBase
admin
2023-01-23 17:23:03
0

在向Hbase中写入数据时,常见的写入方法有使用HBase API,Mapreduce批量导入数据,使用这些方式带入数据时,一条数据写入到HBase数据库中的大致流程如图。
使用BulkLoad从HDFS批量导入数据到HBase

数据发出后首先写入到雨鞋日志WAl中,写入到预写日志中之后,随后写入到内存MemStore中,最后在Flush到Hfile中。这样写数据的方式不会导致数据的丢失,并且道正数据的有序性,但是当遇到大量的数据写入时,写入的速度就难以保证。所以,介绍一种性能更高的写入方式BulkLoad。

使用BulkLoad批量写入数据主要分为两部分:
一、使用HFileOutputFormat2通过自己编写的MapReduce作业将HFile写入到HDFS目录,由于写入到HBase中的数据是按照顺序排序的,HFileOutputFormat2中的configureIncrementalLoad()可以完成所需的配置。
二、将Hfile从HDFS移动到HBase表中,大致过程如图
使用BulkLoad从HDFS批量导入数据到HBase

实例代码pom依赖:


            org.apache.hbase
            hbase-server
            1.4.0
        

        
            org.apache.hadoop
            hadoop-client
            2.6.4
        

        
            org.apache.hbase
            hbase-client
            0.99.2
        
package com.yangshou;

import org.apache.hadoop.hbase.client.Put;
import org.apache.hadoop.hbase.io.ImmutableBytesWritable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapreduce.Mapper;

import java.io.IOException;

public class BulkLoadMapper extends Mapper {
    @Override
    protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException {
        //读取文件中的每一条数据,以序号作为行键
        String line = value.toString();
        //将数据进行切分
        //切分后数组中的元素分别为:序号,用户id,商品id,用户行为,商品分类,时间,地址
        String[] str = line.split(" ");
        String id = str[0];
        String user_id = str[1];
        String item_id = str[2];
        String behavior = str[3];
        String item_type = str[4];
        String time = str[5];
        String address = "156";
        //拼接rowkey和put
        ImmutableBytesWritable rowkry = new ImmutableBytesWritable(id.getBytes());
        Put put = new Put(id.getBytes());
        put.add("info".getBytes(),"user_id".getBytes(),user_id.getBytes());
        put.add("info".getBytes(),"item_id".getBytes(),item_id.getBytes());
        put.add("info".getBytes(),"behavior".getBytes(),behavior.getBytes());
        put.add("info".getBytes(),"item_type".getBytes(),item_type.getBytes());
        put.add("info".getBytes(),"time".getBytes(),time.getBytes());
        put.add("info".getBytes(),"address".getBytes(),address.getBytes());
        //将数据写出
        context.write(rowkry,put);
    }
}
package com.yangshou;

import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.Path;
import org.apache.hadoop.hbase.HBaseConfiguration;
import org.apache.hadoop.hbase.TableName;
import org.apache.hadoop.hbase.client.*;
import org.apache.hadoop.hbase.io.ImmutableBytesWritable;
import org.apache.hadoop.hbase.mapreduce.HFileOutputFormat2;
import org.apache.hadoop.hbase.mapreduce.LoadIncrementalHFiles;
import org.apache.hadoop.mapreduce.Job;
import org.apache.hadoop.mapreduce.lib.input.FileInputFormat;
import org.apache.hadoop.mapreduce.lib.input.TextInputFormat;
import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat;

public class BulkLoadDriver  {
    public static void main(String[] args) throws Exception {
        //获取Hbase配置
        Configuration conf = HBaseConfiguration.create();
        Connection conn = ConnectionFactory.createConnection(conf);
        Table table = conn.getTable(TableName.valueOf("BulkLoadDemo"));
        Admin admin = conn.getAdmin();

        //设置job
        Job job = Job.getInstance(conf,"BulkLoad");
        job.setJarByClass(BulkLoadDriver.class);
        job.setMapperClass(BulkLoadMapper.class);
        job.setMapOutputKeyClass(ImmutableBytesWritable.class);
        job.setMapOutputValueClass(Put.class);

        //设置文件的输入输出路径
        job.setInputFormatClass(TextInputFormat.class);
        job.setOutputFormatClass(HFileOutputFormat2.class);
        FileInputFormat.setInputPaths(job,new Path("hdfs://hadoopalone:9000/tmp/000000_0"));
        FileOutputFormat.setOutputPath(job,new Path("hdfs://hadoopalone:9000/demo1"));

        //将数据加载到Hbase表中
        HFileOutputFormat2.configureIncrementalLoad(job,table,conn.getRegionLocator(TableName.valueOf("BulkLoadDemo")));
        if(job.waitForCompletion(true)){
            LoadIncrementalHFiles load = new LoadIncrementalHFiles(conf);
            load.doBulkLoad(new Path("hdfs://hadoopalone:9000/demo1"),admin,table,conn.getRegionLocator(TableName.valueOf("BulkLoadDemo")));

        }

    }
}

实例数据

44979   100640791   134060896   1   5271    2014-12-09  天津市
44980   100640791   96243605    1   13729   2014-12-02  新疆

在Hbase shell 中创建表

create 'BulkLoadDemo','info'

打包后执行
```hadoop jar BulkLoadDemo-1.0-SNAPSHOT.jar com.yangshou.BulkLoadDriver

注意:在执行hadoop jar之前应该先将Hbase中的相关包加载过来

export HADOOP_CLASSPATH=$HBASE_HOME/lib/*

相关内容

热门资讯

探秘沧海之“晶”:海水制盐的原... 自古以来,海就是巨大的宝藏,尤其它那孕育“白色精灵”——食盐的魔力,更是人类文明的基石之一。海水如何...
日媒播出专题片揭露731部队罪... 当地时间21日夜间,日本广播协会播出侵华日军731部队相关专题片,题为《被隐藏的部队:沉默的战后》。...
苏州正式设立人工智能发展办公室 新京报讯 据“苏州发布”微信公众号消息,近日,苏州市正式设立人工智能发展办公室,进一步统筹苏州市人工...
警方通报“16岁高中生被刺死”... 7月22日,内蒙古乌拉特前旗公安局针对“16岁高中生被刺死,事发地疑为涉黄场所”发布警情通报,通报关...
达观数据陈运文:智能体竞争,正... 随着AI智能体加速进入产业应用,市场规模持续增长,但企业关注的重点正在发生变化。 中商产业研究院发布...
磁力混合器突破3D生物打印细胞... 3D生物打印是生物工程领域的重要技术,通过将活细胞混入柔性水凝胶(即"生物墨水")来打印活体组织,广...
王兴兴:做好准备迎接即将突破的... 王兴兴在开幕式上发表演讲。 王兴兴在演讲中提到,宇树科技成立于2016年,经过近10年努力,已经成为...
不排除提告!屡遭苏巧慧阵营抹黑... 海峡导报综合报道 国民党新北市长参选人李四川的板桥房产问题屡遭绿营质疑,对手苏巧慧阵营再度指控李四川...
俄外长:将在马尼拉同鲁比奥会晤... 据凤凰卫视报道,正在菲律宾出席东盟外长会的俄罗斯外长拉夫罗夫7月22日表示,他将于23日上午在马尼拉...
小鸭洗衣机皮带修理方法 小鸭洗衣机皮带在长时间使用过程中,可能会出现磨损或者老化,导致洗衣机无法转动或者转速变慢。这时候,我...