【发布时间】:2015-08-26 05:14:00
【问题描述】:
假设如果直接从 HDFS 拉取数据而不是使用 HBase API 可以更快地访问数据,我们正在尝试基于来自 HBase 的表快照构建 RDD。
所以,我有一个名为“dm_test_snap”的快照。我似乎能够让大部分配置工作正常工作,但我的 RDD 为空(尽管快照本身中有数据)。
我很难找到任何人使用 Spark 对 HBase 快照进行离线分析的示例,但我不敢相信只有我一个人在尝试让这项工作正常进行。非常感谢任何帮助或建议。
这是我的代码的 sn-p:
object TestSnap {
def main(args: Array[String]) {
val config = ConfigFactory.load()
val hbaseRootDir = config.getString("hbase.rootdir")
val sparkConf = new SparkConf()
.setAppName("testnsnap")
.setMaster(config.getString("spark.app.master"))
.setJars(SparkContext.jarOfObject(this))
.set("spark.executor.memory", "2g")
.set("spark.default.parallelism", "160")
val sc = new SparkContext(sparkConf)
println("Creating hbase configuration")
val conf = HBaseConfiguration.create()
conf.set("hbase.rootdir", hbaseRootDir)
conf.set("hbase.zookeeper.quorum", config.getString("hbase.zookeeper.quorum"))
conf.set("zookeeper.session.timeout", config.getString("zookeeper.session.timeout"))
conf.set("hbase.TableSnapshotInputFormat.snapshot.name", "dm_test_snap")
val scan = new Scan
val job = Job.getInstance(conf)
TableSnapshotInputFormat.setInput(job, "dm_test_snap",
new Path("hdfs://nameservice1/tmp"))
val hBaseRDD = sc.newAPIHadoopRDD(conf, classOf[TableSnapshotInputFormat],
classOf[org.apache.hadoop.hbase.io.ImmutableBytesWritable],
classOf[org.apache.hadoop.hbase.client.Result])
hBaseRDD.count()
System.exit(0)
}
}
更新以包含解决方案 诀窍是,正如下面提到的@Holden,conf 没有通过。为了解决这个问题,我可以通过将 newAPIHadoopRDD 的调用更改为:
val hBaseRDD = sc.newAPIHadoopRDD(job.getConfiguration, classOf[TableSnapshotInputFormat],
classOf[org.apache.hadoop.hbase.io.ImmutableBytesWritable],
classOf[org.apache.hadoop.hbase.client.Result])
@victor 的回答也强调了第二个问题,即我没有通过扫描。为了解决这个问题,我添加了这一行和方法:
conf.set(TableInputFormat.SCAN, convertScanToString(scan))
def convertScanToString(scan : Scan) = {
val proto = ProtobufUtil.toScan(scan);
Base64.encodeBytes(proto.toByteArray());
}
这也让我从 conf.set 命令中拉出这一行:
conf.set("hbase.TableSnapshotInputFormat.snapshot.name", "dm_test_snap")
*注意:这是针对 CDH5.0 上的 HBase 版本 0.96.1.1
便于参考的最终完整代码:
object TestSnap {
def main(args: Array[String]) {
val config = ConfigFactory.load()
val hbaseRootDir = config.getString("hbase.rootdir")
val sparkConf = new SparkConf()
.setAppName("testnsnap")
.setMaster(config.getString("spark.app.master"))
.setJars(SparkContext.jarOfObject(this))
.set("spark.executor.memory", "2g")
.set("spark.default.parallelism", "160")
val sc = new SparkContext(sparkConf)
println("Creating hbase configuration")
val conf = HBaseConfiguration.create()
conf.set("hbase.rootdir", hbaseRootDir)
conf.set("hbase.zookeeper.quorum", config.getString("hbase.zookeeper.quorum"))
conf.set("zookeeper.session.timeout", config.getString("zookeeper.session.timeout"))
val scan = new Scan
conf.set(TableInputFormat.SCAN, convertScanToString(scan))
val job = Job.getInstance(conf)
TableSnapshotInputFormat.setInput(job, "dm_test_snap",
new Path("hdfs://nameservice1/tmp"))
val hBaseRDD = sc.newAPIHadoopRDD(job.getConfiguration, classOf[TableSnapshotInputFormat],
classOf[org.apache.hadoop.hbase.io.ImmutableBytesWritable],
classOf[org.apache.hadoop.hbase.client.Result])
hBaseRDD.count()
System.exit(0)
}
def convertScanToString(scan : Scan) = {
val proto = ProtobufUtil.toScan(scan);
Base64.encodeBytes(proto.toByteArray());
}
}
【问题讨论】:
-
我明白,使用快照而不是实际的 hbase 表的唯一原因是加快进程。但是,必须考虑,当您使用 Hbase 表时,RDD 是从哪里读取的。喜欢它是 HLog 文件还是任何其他文件。一旦确认了该方面,快照和实际表在上述方面是相似的。我们在与 hbase 的外部框架集成方面遇到了类似的问题。如果我们采用传统方法,一切都很好。任何新的节省时间的东西,框架都有一些限制。
-
我希望通过 HDFS 通过快照直接访问 HFiles,而收益将是直接从磁盘将数据流式传输到 RDD,绕过对 HBase 的任何调用。
-
快照包含对拍摄快照时表中文件的引用。在快照操作期间不会制作数据副本,但可能会在触发压缩或删除时制作副本。所以 newAPIHadoopRDD() 方法需要有额外的逻辑来从快照中获取实际的 HFile,而不是定期查找 hadoop/hbase 文件。需要在 RDD 级别确认此行为
-
拥有日志可能会有所帮助,这样我们就可以验证设置的信息是否一直向下传递。
-
我将此代码用于类似的用例,而不是我正在使用的“工作”对象
TableSnapshotInputFormatImpl.setInput(config ,snapShotName,path)这工作正常,但数据提取在相同范围“R/0”时非常慢"-"R/1",这种方法需要 4 小时才能提取 85GB 的数据,而查询具有相同范围的表的正常作业在 10 分钟内完成,不知道可能是什么问题。 Spark DAG 和两个工作的“解释”计划完全相同。
标签: scala hadoop apache-spark hbase