【问题标题】:Spark write only to one hbase region serverSpark 仅写入一个 hbase 区域服务器
【发布时间】:2017-02-03 18:23:00
【问题描述】:
import org.apache.hadoop.hbase.mapreduce.TableOutputFormat
import org.apache.hadoop.hbase.mapreduce.TableInputFormat
import org.apache.hadoop.mapreduce.Job
import org.apache.hadoop.hbase.io.ImmutableBytesWritable
import org.apache.spark.rdd.PairRDDFunctions

def bulkWriteToHBase(sparkSession: SparkSession, sparkContext: SparkContext, jobContext: Map[String, String], sinkTableName: String, outRDD: RDD[(ImmutableBytesWritable, Put)]): Unit = {
val hConf = HBaseConfiguration.create()
hConf.set("hbase.zookeeper.quorum", jobContext("hbase.zookeeper.quorum"))
hConf.set("zookeeper.znode.parent", jobContext("zookeeper.znode.parent"))
hConf.set(TableInputFormat.INPUT_TABLE, sinkTableName)

val hJob = Job.getInstance(hConf)
hJob.getConfiguration().set(TableOutputFormat.OUTPUT_TABLE, sinkTableName)
hJob.setOutputFormatClass(classOf[TableOutputFormat[ImmutableBytesWritable]]) 

outRDD.saveAsNewAPIHadoopDataset(hJob.getConfiguration())
}

我使用这个hbase批量插入发现的是,每次spark只会从hbase写入一个单一的区域服务器,这成为了瓶颈。

但是,当我使用几乎相同的方法但从 hbase 读取时,它使用多个执行器进行并行读取。

def bulkReadFromHBase(sparkSession: SparkSession, sparkContext: SparkContext, jobContext: Map[String, String], sourceTableName: String) = {
val hConf = HBaseConfiguration.create()
hConf.set("hbase.zookeeper.quorum", jobContext("hbase.zookeeper.quorum"))
hConf.set("zookeeper.znode.parent", jobContext("zookeeper.znode.parent"))
hConf.set(TableInputFormat.INPUT_TABLE, sourceTableName)

val inputRDD = sparkContext.newAPIHadoopRDD(hConf, classOf[TableInputFormat], classOf[ImmutableBytesWritable], classOf[Result])
inputRDD
}

谁能解释一下为什么会发生这种情况?或者我有 对 spark-hbase 批量 I/O 使用了错误的方式?

【问题讨论】:

    标签: apache-spark hadoop hbase rdd


    【解决方案1】:

    问题:我对 spark-hbase 批量 I/O 使用了错误的方式?

    虽然您的方法不对,但您需要事先预先分割区域并创建带有预先分割区域的表格。

    例如create 'test_table', 'f1', SPLITS=> ['1', '2', '3', '4', '5', '6', '7', '8', '9']

    上表占据9个区域..

    设计好的 rowkey 以 1-9 开头

    您可以使用如下所示的番石榴杂音哈希。

    import com.google.common.hash.HashCode;
    import com.google.common.hash.HashFunction;
    import com.google.common.hash.Hashing;
    
    /**
         * getMurmurHash.
         * 
         * @param content
         * @return HashCode
         */
        public static HashCode getMurmurHash(String content) {
            final HashFunction hf = Hashing.murmur3_128();
            final HashCode hc = hf.newHasher().putString(content, Charsets.UTF_8).hash();
            return hc;
        }
    
    final long hash = getMurmur128Hash(Bytes.toString(yourrowkey as string)).asLong();
                final int prefix = Math.abs((int) hash % 9);
    

    现在将此前缀附加到您的行键

    例如

    1rowkey1 // 将进入第一个区域
    2rowkey2 // 将进入 第二区
    3rowkey3 // 将进入第三个区域 ... 9rowkey9 // 将进入第九区

    如果您正在进行预拆分,并且想要手动管理区域拆分,您还可以通过将 hbase.hregion.max.filesize 设置为较大的数字并将拆分策略设置为 ConstantSizeRegionSplitPolicy 来禁用区域拆分。但是,您应该使用像 100GB 这样的安全值,这样区域就不会超出区域服务器的能力。您可以考虑禁用自动拆分并依赖于预拆分的初始区域集,例如,如果您使用统一哈希作为键前缀,并且您可以确保每个区域的读/写负载区域及其大小在表中的区域之间是统一的

    1) 请确保在将数据加载到 hbase 表之前可以预先拆分表 2) 使用 murmurhash 或其他一些散列技术设计好的行键,如下所述。以确保跨地区的均匀分布。
    也可以看看http://hortonworks.com/blog/apache-hbase-region-splitting-and-merging/

    问题:谁能解释一下为什么会发生这种情况?

    原因非常明显和简单由于该表的行键不佳,数据的热点是一个特定的原因...

    考虑 java 中的 hashmap,它的元素的 hashcode 为 1234。那么它将填充一个桶中的所有元素,不是吗?如果 hashmap 元素分布在不同的好 hashcode 中,那么它会将元素放在不同的桶中。 hbase 也是如此。在这里,您的哈希码就像您的 rowkey...

    还有更多,

    如果我已经有一张桌子并且我想分割区域会发生什么 跨越...

    RegionSplitter 类提供了几个实用程序来帮助选择手动分割区域而不是让 HBase 自动处理的开发人员管理生命周期。

    最有用的实用程序是:

    • 创建具有指定数量的预分割区域的表
    • 对现有表上的所有区域执行滚动拆分

    例子:

    $ hbase org.apache.hadoop.hbase.util.RegionSplitter test_table HexStringSplit -c 10 -f f1
    

    其中-c 10,指定请求的区域数为10,-f指定表中你想要的列族,用“:”隔开。该工具将创建一个名为“test_table”的表,其中包含 10 个区域:

    13/01/18 18:49:32 DEBUG hbase.HRegionInfo: Current INFO from scan results = {NAME => 'test_table,,1358563771069.acc1ad1b7962564fc3a43e5907e8db33.', STARTKEY => '', ENDKEY => '19999999', ENCODED => acc1ad1b7962564fc3a43e5907e8db33,}
    13/01/18 18:49:32 DEBUG hbase.HRegionInfo: Current INFO from scan results = {NAME => 'test_table,19999999,1358563771096.37ec12df6bd0078f5573565af415c91b.', STARTKEY => '19999999', ENDKEY => '33333332', ENCODED => 37ec12df6bd0078f5573565af415c91b,}
    ...
    

    正如评论中所讨论的,您发现我在写入 hbase 之前的最终 RDD 只有 1 个分区!这表明有 只有一个执行者持有整个数据......我仍在尝试 找出原因。

    另外,检查

    spark.default.parallelism 默认为所有内核上的所有内核数 机器。并行化 api 没有父 RDD 来确定 分区数,所以它使用spark.default.parallelism。

    所以你可以通过重新分区来增加分区。

    注意:我观察到,在 Mapreduce 中,区域的分区数/输入拆分 = 启动的映射器数。同样,在您的情况下,数据加载到一个特定区域的情况可能相同,这就是为什么一个执行程序发射。请同时验证

    【讨论】:

    • 感谢您的回答,我发现我在写入 hbase 之前的最终 RDD 只有 1 个分区!这表明只有一个执行者持有整个数据......我仍在试图找出原因。在这种情况下,预拆分甚至无济于事,因为只有一个执行程序在运行并将数据写入 hbase。
    • 是的,这也应该小心!检查默认并行属性,因为您没有重新分区。 spark.default.parallelism 默认为所有机器上所有核心的数量。 parallelize api 没有父 RDD 来确定分区数,所以它使用 spark.default.parallelism。
    • 我在批量写入 hbase 之前尝试了重新分区,并且成功了,所以现在我看到数据被同时写入多个区域。但是我尝试设置 spark.default.parallelism=64 并没有做任何事情。所以现在我正试图找到数据在哪里折叠到一个分区中。感谢您的帮助,这很有帮助。
    • 太棒了..如果您还没有投票,请注意投票,如果您还可以!谢谢!
    【解决方案2】:

    虽然您没有提供示例数据或足够的解释,但这主要不是由于您的代码或配置。 由于非最佳行键设计,它正在发生。 您正在编写的数据的键(hbase rowkey)结构不正确(可能是单调增加或其他)。因此,正在写入其中一个区域。您可以通过各种方式防止这种情况(rowkey 设计的各种推荐做法,例如盐渍、反转和其他技术)。 供参考,您可以查看http://hbase.apache.org/book.html#rowkey.design

    如果您想知道写入是针对所有区域并行完成还是一个一个完成(从问题中不清楚),请查看: http://hbase.apache.org/book.html#_bulk_load.

    【讨论】:

    • 感谢您的帮助,我确实发现我的 RDD 通过一系列数据帧操作最终进入了 1 个分区......虽然我不知道哪个操作触发了分区崩溃到一。我确实从 hbase 读取到 64 个分区,我可以通过查看 sparkUI 阶段来证明这一点。
    • 即使您的 rdd 仅位于 spark 的一个分区中,它也会根据密钥结构写入不同的 hbase 区域。我想您是在问为什么 spark 不以并行方式写入 hbase (一个或多个区域,没关系)使用多个执行器。如果是这样,请告诉我。
    • 是的,你是对的。它是根据region splitkeys写入多个region的,所以只是只有一个executor在写出,影响IO性能。
    猜你喜欢
    • 2013-12-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-07-23
    • 1970-01-01
    相关资源
    最近更新 更多