這篇文章將為大家詳細(xì)講解有關(guān)Spark如何批量存取HBase ,小編覺得挺實用的,因此分享給大家做個參考,希望大家閱讀完這篇文章后可以有所收獲。
讓客戶滿意是我們工作的目標(biāo),不斷超越客戶的期望值來自于我們對這個行業(yè)的熱愛。我們立志把好的技術(shù)通過有效、簡單的方式提供給客戶,將通過不懈努力成為客戶在信息化領(lǐng)域值得信任、有價值的長期合作伙伴,公司提供的服務(wù)項目有:域名與空間、網(wǎng)頁空間、營銷軟件、網(wǎng)站建設(shè)、龍湖網(wǎng)站維護、網(wǎng)站推廣。
FileAna.scala
object FileAna {
// val conf: Configuration = HBaseConfiguration.create()
val hdfsPath = "hdfs://master:9000"
val hdfs = FileSystem.get(new URI(hdfsPath), new Configuration())
def main(args: Array[String]) {
val conf = new SparkConf().setAppName("FileAna").setMaster("spark://master:7077").
set("spark.driver.host", "192.168.1.127").
setJars(List("/home/pang/woozoomws/spark-service.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-common-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-client-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-protocol-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/htrace-core-3.1.0-incubating.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-server-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/metrics-core-2.2.0.jar"))
val sc = new SparkContext(conf)
val rdd = sc.textFile("hdfs://master:9000/woozoom/msgfile.txt")
val rdd2 = rdd.map(x => convertToHbase(anaMavlink(x)))
val hbaseConf = HBaseConfiguration.create()
hbaseConf.addResource("/home/hadoop/software/hbase-1.2.2/conf/hbase-site.xml");
val jobConf = new JobConf(hbaseConf, this.getClass)
jobConf.setOutputFormat(classOf[TableOutputFormat])
jobConf.set(TableOutputFormat.OUTPUT_TABLE, "MissionItem")
rdd2.saveAsHadoopDataset(jobConf)
sc.stop()
}
def convertScanToString(scan: Scan) = {
val proto = ProtobufUtil.toScan(scan)
Base64.encodeBytes(proto.toByteArray)
}
def convertToHbase(msg: MAVLinkMessage) = {
val p = new Put(Bytes.toBytes(UUID.randomUUID().toString()))
if (msg.isInstanceOf[msg_mission_item]) {
val missionItem = msg.asInstanceOf[msg_mission_item]
p.addColumn(Bytes.toBytes("data"), Bytes.toBytes("x"), Bytes.toBytes(missionItem.x))
p.addColumn(Bytes.toBytes("data"), Bytes.toBytes("y"), Bytes.toBytes(missionItem.y))
p.addColumn(Bytes.toBytes("data"), Bytes.toBytes("z"), Bytes.toBytes(missionItem.z))
}
(new ImmutableBytesWritable, p)
}
val anaMavlink = (str: String) => {
val bytes = ByteAndHex.hexStringToBytes(str)
QuickParser.parse(bytes).unpack()
}
}
ReadHBase.scala
object ReadHBase {
// val conf: Configuration = HBaseConfiguration.create()
val hdfsPath = "hdfs://master:9000"
val hdfs = FileSystem.get(new URI(hdfsPath), new Configuration())
def main(args: Array[String]) {
val conf = new SparkConf().setAppName("FileAna").setMaster("spark://master:7077").
set("spark.driver.host", "192.168.1.127").
setJars(List("/home/pang/woozoomws/spark-service.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-common-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-client-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-protocol-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/htrace-core-3.1.0-incubating.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/hbase-server-1.2.2.jar",
"/home/pang/woozoomws/spark-service/lib/hbase/metrics-core-2.2.0.jar"))
val sc = new SparkContext(conf)
val hbaseConf = HBaseConfiguration.create()
hbaseConf.addResource("/home/hadoop/software/hbase-1.2.2/conf/hbase-site.xml");
hbaseConf.set(TableInputFormat.INPUT_TABLE, "MissionItem")
val scan = new Scan()
hbaseConf.set(TableInputFormat.SCAN, convertScanToString(scan))
val readRDD = sc.newAPIHadoopRDD(hbaseConf, classOf[TableInputFormat],
classOf[org.apache.hadoop.hbase.io.ImmutableBytesWritable],
classOf[org.apache.hadoop.hbase.client.Result])
val count = readRDD.count()
println("Mission Item Count:" + count)
sc.stop()
}
def convertScanToString(scan: Scan) = {
val proto = ProtobufUtil.toScan(scan)
Base64.encodeBytes(proto.toByteArray)
}
}
關(guān)于“Spark如何批量存取HBase ”這篇文章就分享到這里了,希望以上內(nèi)容可以對大家有一定的幫助,使各位可以學(xué)到更多知識,如果覺得文章不錯,請把它分享出去讓更多的人看到。