【问题标题】:How to do parallel reading from Hbase using list of RowKeys in Spark如何使用 Spark 中的 RowKeys 列表从 Hbase 进行并行读取
【发布时间】:2018-09-01 00:31:19
【问题描述】:

我有 650 万行键,想从 spark-job 中的 hbase 检索数据。如何从 hbase 并行检索结果?

我认为这个 sn-p 不会在执行程序上运行。

List<Get> listOFGets = new ArrayList<Get>();
Result[] results = Htable.get(listOFGets);

【问题讨论】:

标签: apache-spark hbase


【解决方案1】:

我通常使用.newAPIHadoopRDD() 方法进行hbase 扫描。请注意,这是 scala 与 java api 的非常丑陋的组合。您可以传入任意行键列表(空列表返回表中的所有记录)。如果您的行键不是长编码的,那么您可能需要稍微修改一下代码。

def hbaseScan(ids: List[Long]): Dataset[Result] = {
  val ranges = ListBuffer[MultiRowRangeFilter.RowRange]()
  //converts each id (Long) into a one element RowRange
  //(id gets implicitly get converted to byte[])
  ids.foreach(i => {
    ranges += new MultiRowRangeFilter.RowRange(i, true, i + 1, false)
  })

  val scan = new Scan()
  scan.setCaching(1000) /* fetch 1000 records in each trip to hbase */
  scan.setCacheBlocks(false) /* don't waste hbase cache space, since we are scanning whole table

  if (ranges.nonEmpty) {
    //The list of RowRanges is sorted and merged into a single scan filter
    scan.setFilter(new MultiRowRangeFilter(MultiRowRangeFilter.sortAndMerge(ranges.asJava)))
  }

  val conf = HBaseConfiguration.create()
  conf.set(TableInputFormat.INPUT_TABLE, HBASE_TABLE /*set your table name here*/)
  conf.set(TableInputFormat.SCAN, scan)

  spark.sparkContext.newAPIHadoopRDD(
    conf,
    classOf[TableInputFormat],
    classOf[ImmutableBytesWritable],
    classOf[Result]
  ).toDF("result").as[Result]
}

这将返回一个Dataset[Result],其分区数与扫描表中的区域数相同。抱歉,我没有任何等效的 java 代码可以分享。

编辑:解决不正确的评论方式

我应该先说这种方法在读取整个 hbase 表或少量任意行键时效果最好。我的用例正是在做这两个方面,因为我总是一次查询 1000 个行键,或者整个表,中间什么都没有。

如果您的任意行键数量很大,则在MultiRowRangeFilter.sortAndMerge() 步骤中将有一个核心挂起。此方法可以扩展为在创建Filter 进行扫描之前并行化排序和合并键列表到键范围的过程。在排序和合并之后,这种方法确实可以在尽可能多的分区上并行,如果您有许多连续的行键范围,甚至可以减少到 hbase 的往返次数。

很难说这个过程是否比在集群中散布随机数据更有效,因为它完全取决于许多因素:记录大小、表大小、行键范围等。我相信对于许多人来说用例这种方法会更有效,但显然并非适用于所有用例。

【讨论】:

  • 这不是正确的方法。甚至您可以在 spark 上而不是在 java 上的 htable.get 上监控这个逻辑。
  • “不正确”是一个不正确的评估。这种方法在从 hbase 检索数据时确实是并行的,但我承认它有一个弱点,可以改进,如编辑中所述。
【解决方案2】:

您可以通过创建RDD(以适合您的方式)然后调用JavaHBaseContext 对象的方法.foreachPartition 来对执行程序运行并行扫描。这样,HBaseContext 会将 Connection 实例详细信息传递给函数类的每个实例,并且在该函数内部,您可以通过获取表等方式进行扫描。

如何创建 RDD 来适应这种情况,以及它应该有多少元素(只要确保它具有与并行扫描一样多的分区),这取决于您。根据我的经验,您可以在 spark 上运行多少并发任务,可以执行多少并行扫描(当然取决于您的 HBase 集群强度)。

Java 代码可能如下所示:

在主人身上:

JavaHBaseContext hBaseContext = new JavaHBaseContext(sparkContext, HBaseConfig);
JavaRDD<blah> myRDD = ... (create an RDD with a number of elements)
hBaseContext.foreachPartition(myRDD,  new MyParFunction());

您的函数类将如下所示:

class MyParFunction implements VoidFunction<Tuple2<Iterator<blah>, Connection>>
{
@Override
    public void call(Tuple2<Iterator<blah>, Connection> t) throws Exception
    {
// Do your scan here, since you have the Connection object t
}
}

这应该在所有执行器上并行运行扫描

【讨论】:

  • 太棒了。接受答案,以便其他人将来知道这是有效的。
猜你喜欢
  • 1970-01-01
  • 2015-01-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-10-01
  • 1970-01-01
相关资源
最近更新 更多