温馨提示×

温馨提示×

您好,登录后才能下订单哦!

密码登录×
登录注册×
其他方式登录
点击 登录注册 即表示同意《亿速云用户服务条款》

怎么用Spark读取HBASE数据

发布时间:2021-08-18 11:02:33 来源:亿速云 阅读:808 作者:chen 栏目:web开发

这篇文章主要讲解了“怎么用Spark读取HBASE数据”,文中的讲解内容简单清晰,易于学习与理解,下面请大家跟着小编的思路慢慢深入,一起来研究和学习“怎么用Spark读取HBASE数据”吧!

scala访问HBASE通常2种方式,一种是使用SPARK方式读取HBASE数据直接转换成RDD, 一种采用和JAVA类似的方式,通过HTable操作HBASE,数据获取之后再自己进行处理。 这2种方式区别应该是RDD是跑在多节点通过从HBASE获取数据,而采用HTable的方式,应该是串行了,仅仅是HBASE层面是分布式而已。

1. 转换为RDD
package com.isesol.spark
import org.apache.spark.SparkContext
import org.apache.spark.SparkConf
import org.apache.hadoop.hbase.HBaseConfiguration
import org.apache.hadoop.mapred.JobConf
import org.apache.hadoop.hbase.mapred.TableOutputFormat
import org.apache.hadoop.hbase.util.Bytes
import org.apache.spark.rdd.RDD.rddToPairRDDFunctions
import org.apache.hadoop.hbase.client.Put
import org.apache.hadoop.hbase.io.ImmutableBytesWritable
import org.apache.hadoop.hbase.spark._
import org.apache.hadoop.hbase.client.Scan
import org.apache.hadoop.hbase.TableName
import org.apache.hadoop.hbase.filter._
import org.apache.hadoop.hbase.filter.CompareFilter.CompareOp


object hbasescan {


  def main(args: Array[String]) {
    val conf = new SparkConf()
    conf.setMaster("local").setAppName("this is for spark SQL")
    //conf.setSparkHome("d:\\spark_home")
    val hbaseconf = HBaseConfiguration.create()
    hbaseconf.set("hbase.zookeeper.quorum", "datanode01.isesol.com,datanode02.isesol.com,datanode03.isesol.com,datanode04.isesol.com,cmserver.isesol.com")
    hbaseconf.set("hbase.zookeeper.property.clientPort", "2181")
    hbaseconf.set("maxSessionTimeout", "6")
    val sc = new SparkContext(conf)
    try {
      println("start to read from hbase")
      val hbaseContext = new HBaseContext(sc, hbaseconf)
      val scan = new Scan()
      scan.setMaxVersions()
      //scan.setRowPrefixFilter(Bytes.toBytes("i51530048-1007-9223370552914159518"))
      scan.setCaching(100)
      val filter = new SingleColumnValueFilter(Bytes.toBytes("cf"), Bytes.toBytes("age"), CompareOp.LESS, Bytes.toBytes("1"));
      scan.setFilter(filter)
      val hbaserdd = hbaseContext.hbaseRDD(TableName.valueOf("bank"), scan)
      hbaserdd.cache()
      println(hbaserdd.count())
    } catch {
      case ex: Exception => println("can not connect hbase")
    }
  }
}

2. 采用 HTable方式处理

    val htable = new HTable(hbaseconf, "t_device_fault_statistics")
    val scan1 = new Scan()
    scan1.setCaching(3*1024*1024)
    val scaner = htable.getScanner(scan1)
    
    while(scaner.iterator().hasNext()){
       val result = scaner.next()
       if(result.eq(null)){
       } else {
         println(Bytes.toString(result.getRow) + "\t" + Bytes.toString(result.getValue("cf".getBytes, "fault_level2_name".getBytes)))
       }
    }
    scaner.close()
    htable.close()

感谢各位的阅读,以上就是“怎么用Spark读取HBASE数据”的内容了,经过本文的学习后,相信大家对怎么用Spark读取HBASE数据这一问题有了更深刻的体会,具体使用情况还需要大家实践验证。这里是亿速云,小编将为大家推送更多相关知识点的文章,欢迎关注!

向AI问一下细节

免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。

AI