【问题标题】:How to process a range of hbase rows using spark?如何使用 spark 处理一系列 hbase 行?
【发布时间】:2014-10-01 02:36:36
【问题描述】:

我正在尝试使用 HBase 作为 spark 的数据源。所以第一步是从 HBase 表创建一个 RDD。由于 Spark 使用 hadoop 输入格式,我可以通过创建 rdd http://www.vidyasource.com/blog/Programming/Scala/Java/Data/Hadoop/Analytics/2014/01/25/lighting-a-spark-with-hbase 找到使用所有行的方法但是我们如何为范围扫描创建 RDD?

欢迎所有建议。

【问题讨论】:

    标签: java hadoop bigdata apache-spark


    【解决方案1】:

    这是在 Spark 中使用 Scan 的示例:

    import java.io.{DataOutputStream, ByteArrayOutputStream}
    import java.lang.String
    import org.apache.hadoop.hbase.client.Scan
    import org.apache.hadoop.hbase.HBaseConfiguration
    import org.apache.hadoop.hbase.io.ImmutableBytesWritable
    import org.apache.hadoop.hbase.client.Result
    import org.apache.hadoop.hbase.mapreduce.TableInputFormat
    import org.apache.hadoop.hbase.util.Base64
    
    def convertScanToString(scan: Scan): String = {
      val out: ByteArrayOutputStream = new ByteArrayOutputStream
      val dos: DataOutputStream = new DataOutputStream(out)
      scan.write(dos)
      Base64.encodeBytes(out.toByteArray)
    }
    
    val conf = HBaseConfiguration.create()
    val scan = new Scan()
    scan.setCaching(500)
    scan.setCacheBlocks(false)
    conf.set(TableInputFormat.INPUT_TABLE, "table_name")
    conf.set(TableInputFormat.SCAN, convertScanToString(scan))
    val rdd = sc.newAPIHadoopRDD(conf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result])
    rdd.count
    

    您需要将相关库添加到 Spark 类路径并确保它们与您的 Spark 兼容。提示:您可以使用hbase classpath 找到它们。

    【讨论】:

    • 您使用的是什么版本的 HBase?在我扫描的任何版本中都找不到Scan.write 方法。
    • Hbase 现在使用 Protobuf。你可以在这里找到如何实现convertScanToStringgithub.com/apache/hbase/blob/…
    • 这个 setCaching(500) 会在 HBase 中创建 500 行的 rdd 吗?我试过了,它仍然从 Hbase 获取所有数据。
    • 没有。客户端每次将请求 500 行,但仍会获取所有数据。
    • 为了让导入工作,我不得不使用org.apache.hbase:hbase-client:1.1.2 org.apache.hbase:hbase-common:1.1.2 org.apache.hbase:hbase-server:1.1.2
    【解决方案2】:

    您可以在下面设置conf

     val conf = HBaseConfiguration.create()//need to set all param for habse
     conf.set(TableInputFormat.SCAN_ROW_START, "row2");
     conf.set(TableInputFormat.SCAN_ROW_STOP, "stoprowkey");
    

    这将只为那些记录加载 rdd

    【讨论】:

      【解决方案3】:

      这是一个带有TableMapReduceUtil.convertScanToString(Scan scan):的Java示例

      import org.apache.hadoop.conf.Configuration;
      import org.apache.hadoop.hbase.HBaseConfiguration;
      import org.apache.hadoop.hbase.HConstants;
      import org.apache.hadoop.hbase.client.Result;
      import org.apache.hadoop.hbase.client.Scan;
      import org.apache.hadoop.hbase.io.ImmutableBytesWritable;
      import org.apache.hadoop.hbase.mapreduce.TableInputFormat;
      import org.apache.hadoop.hbase.mapreduce.TableMapReduceUtil;
      import org.apache.spark.SparkConf;
      import org.apache.spark.api.java.JavaPairRDD;
      import org.apache.spark.api.java.JavaSparkContext;
      
      import java.io.IOException;
      
      public class HbaseScan {
      
          public static void main(String ... args) throws IOException, InterruptedException {
      
              // Spark conf
              SparkConf sparkConf = new SparkConf().setMaster("local[4]").setAppName("My App");
              JavaSparkContext jsc = new JavaSparkContext(sparkConf);
      
              // Hbase conf
              Configuration conf = HBaseConfiguration.create();
              conf.set(TableInputFormat.INPUT_TABLE, "big_table_name");
      
              // Create scan
              Scan scan = new Scan();
              scan.setCaching(500);
              scan.setCacheBlocks(false);
              scan.setStartRow(Bytes.toBytes("a"));
              scan.setStopRow(Bytes.toBytes("d"));
      
      
              // Submit scan into hbase conf
              conf.set(TableInputFormat.SCAN, TableMapReduceUtil.convertScanToString(scan));
      
              // Get RDD
              JavaPairRDD<ImmutableBytesWritable, Result> source = jsc
                      .newAPIHadoopRDD(conf, TableInputFormat.class,
                              ImmutableBytesWritable.class, Result.class);
      
              // Process RDD
              System.out.println(source.count());
          }
      }
      

      【讨论】:

        猜你喜欢
        • 2016-03-06
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多