【问题标题】:How to connect with Hbase using spark如何使用 spark 连接 Hbase
【发布时间】:2015-08-12 11:46:01
【问题描述】:

我想使用 spark 连接 hbase。 我遇到了一个例外。当我试图从 scala 做同样的事情时,我没有收到这样的错误。

我正在使用 scala 2.10.4,火花:- 1.3.0 CDH 5.4.0。

代码是:

import org.apache.spark.SparkContext
import org.apache.hadoop.hbase.HBaseConfiguration
import org.apache.hadoop.fs.Path
import org.apache.hadoop.hbase.util.Bytes
import org.apache.hadoop.hbase.client.Put
import org.apache.spark.SparkConf
import com.cloudera.spark.hbase.HBaseContext

object HBaseBulkPutExample {
    def main(args: Array[String]) {
        val tableName = "image_table";
        val columnFamily = "Image Data"

        val sparkConf = new SparkConf().setAppName("HBaseBulkPutExample " + tableName + " " + columnFamily)
        val sc = new SparkContext(sparkConf)

        val rdd = sc.parallelize(Array(
            (Bytes.toBytes("1"), Array((Bytes.toBytes(columnFamily), Bytes.toBytes("1"), Bytes.toBytes("1")))),
            (Bytes.toBytes("2"), Array((Bytes.toBytes(columnFamily), Bytes.toBytes("1"), Bytes.toBytes("2")))),
            (Bytes.toBytes("3"), Array((Bytes.toBytes(columnFamily), Bytes.toBytes("1"), Bytes.toBytes("3")))),
            (Bytes.toBytes("4"), Array((Bytes.toBytes(columnFamily), Bytes.toBytes("1"), Bytes.toBytes("4")))),
            (Bytes.toBytes("5"), Array((Bytes.toBytes(columnFamily), Bytes.toBytes("1"), Bytes.toBytes("5"))))
            )
        )

        val conf = HBaseConfiguration.create();
        conf.addResource(new Path("/eds/servers//hbase-1.0.1.1/conf/hbase-site.xml"));

        val hbaseContext = new HBaseContext(sc, conf);
        hbaseContext.bulkPut[(Array[Byte], Array[(Array[Byte], Array[Byte], Array[Byte])])](rdd,
           tableName,
           (putRecord) => {
               val put = new Put(putRecord._1)
               putRecord._2.foreach((putValue) => put.add(putValue._1, putValue._2, putValue._3))
               put
            },
            true);
    }
}

当我创建一个 jar 并执行它时,我收到以下错误:

org.apache.hadoop.hbase.DoNotRetryIOException: java.lang.NoSuchMethodError: org.apache.hadoop.net.NetUtils.getInputStream(Ljava/net/Socket;)Lorg/apache/hadoop/net/SocketInputWrapper;
    at org.apache.hadoop.hbase.ipc.RpcClient$Connection.setupIOstreams(RpcClient.java:928)
    at org.apache.hadoop.hbase.ipc.RpcClient.getConnection(RpcClient.java:1543)
    at org.apache.hadoop.hbase.ipc.RpcClient.call(RpcClient.java:1442)
    at org.apache.hadoop.hbase.ipc.RpcClient.callBlockingMethod(RpcClient.java:1661)
    at org.apache.hadoop.hbase.ipc.RpcClient$BlockingRpcChannelImplementation.callBlockingMethod(RpcClient.java:1719)
    at org.apache.hadoop.hbase.protobuf.generated.ClientProtos$ClientService$BlockingStub.get(ClientProtos.java:30304)
    at org.apache.hadoop.hbase.protobuf.ProtobufUtil.getRowOrBefore(ProtobufUtil.java:1562)
    at org.apache.hadoop.hbase.client.HTable$2.call(HTable.java:711)
    at org.apache.hadoop.hbase.client.HTable$2.call(HTable.java:709)
    at org.apache.hadoop.hbase.client.RpcRetryingCaller.callWithRetries(RpcRetryingCaller.java:114)
    at org.apache.hadoop.hbase.client.HTable.getRowOrBefore(HTable.java:715)
    at org.apache.hadoop.hbase.client.MetaScanner.metaScan(MetaScanner.java:144)
    at org.apache.hadoop.hbase.client.HConnectionManager$HConnectionImplementation.prefetchRegionCache(HConnectionManager.java:1140)

【问题讨论】:

    标签: scala apache-spark hbase


    【解决方案1】:

    由于 spark、hbase 中的库中的版本不匹配问题,我遇到了类似的问题。我升级到 spark 1.4、scala 2.11.6 并使用 hbase 1.0.1.1 - 之后再也没有遇到过这个问题。 Jar 版本不匹配导致此问题,因为 spark jar 期望在 hbase 客户端 jar 中升级方法并失败。

    【讨论】:

    • 我无法添加以下工件,我可以在其中找到 HBaseContext,我使用的是 Hbase 1.1。 org.apache.hbasehbase1.1.1
    【解决方案2】:

    我不确定您的错误的原因是什么。看起来像调用 Method mismatch 与您的环境版本。

    这是一个关于如何使用 spark 连接 Hbase 的示例:

    import spark._
    import org.apache.hadoop.hbase.{HBaseConfiguration, HTableDescriptor}
    import org.apache.hadoop.hbase.client.HBaseAdmin
    import org.apache.hadoop.hbase.mapreduce.TableInputFormat
    import org.apache.hadoop.hbase.client.Result
    import org.apache.hadoop.hbase.io.ImmutableBytesWritable
    import org.apache.hadoop.hbase.mapreduce.TableInputFormat
    
    ...
    
    val conf = HBaseConfiguration.create()
    conf.set(TableInputFormat.INPUT_TABLE, "image_table")
    
     // Initialize 
    val admin = new HBaseAdmin(conf)
    if(!admin.isTableAvailable(input_table)) {
      val tableDesc = new HTableDescriptor("image_table")
      admin.createTable(tableDesc)
    }
    
    val rdd = sc.newAPIHadoopRDD(
        conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result])
    

    Result 类包括各种获取值的方法

    【讨论】:

    • 您好,感谢您的回复。你是对的 Azi,我看到其他人也提出了同样的建议,但我不明白哪种方法是正确的。如果你真的在做这个,你能给我完整的代码吗(带管理员,put,...)。
    • 请参考使用 spark 更新的 Hbase 骨架。仅供参考:管理员是通用的,您可以参考 hbase 文档站点。
    猜你喜欢
    • 2016-11-23
    • 1970-01-01
    • 2017-02-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-11-04
    相关资源
    最近更新 更多