【问题标题】:Insert Spark dataframe into hbase将 Spark 数据帧插入 hbase
【发布时间】:2017-05-22 11:42:36
【问题描述】:

我有一个数据框,我想将它插入到 hbase 中。我关注这个documenation

这就是我的数据框的样子:

 --------------------
|id | name | address |
|--------------------|
|23 |marry |france   |
|--------------------|
|87 |zied  |italie   |
 --------------------

我使用以下代码创建了一个 hbase 表:

val tableName = "two"
val conf = HBaseConfiguration.create()
if(!admin.isTableAvailable(tableName)) {
          print("-----------------------------------------------------------------------------------------------------------")
          val tableDesc = new HTableDescriptor(tableName)
          tableDesc.addFamily(new HColumnDescriptor("z1".getBytes()))
          admin.createTable(tableDesc)
        }else{
          print("Table already exists!!--------------------------------------------------------------------------------------")
        }

现在如何将这个数据帧插入到 hbase 中?

在另一个示例中,我使用以下代码成功插入 hbase:

val myTable = new HTable(conf, tableName)
    for (i <- 0 to 1000) {
      var p = new Put(Bytes.toBytes(""+i))
      p.add("z1".getBytes(), "name".getBytes(), Bytes.toBytes(""+(i*5)))
      p.add("z1".getBytes(), "age".getBytes(), Bytes.toBytes("2017-04-20"))
      p.add("z2".getBytes(), "job".getBytes(), Bytes.toBytes(""+i))
      p.add("z2".getBytes(), "salary".getBytes(), Bytes.toBytes(""+i))
      myTable.put(p)
    }
    myTable.flushCommits()

但是现在我被卡住了,如何将我的数据帧的每条记录插入到我的 hbase 表中。

感谢您的时间和关注

【问题讨论】:

  • 问题不清楚。你在做别的事情。 hbase.apache.org/book.html#_sparksql_dataframes 告诉您定义目录并在 sc.parallelize(data).toDF.write.options 中使用以将 DF 保存到 HBase。
  • 是的,并提到我正在使用该文档。我被困在这里val data = (0 to 255).map { i =&gt; HBaseRecord(i, "extra")} 如何插入我的数据帧的 foreach 记录,而不是从 0 到 255

标签: scala apache-spark dataframe hbase rdd


【解决方案1】:

另一种方法是查看 rdd.saveAsNewAPIHadoopDataset,将数据插入到 hbase 表中。

def main(args: Array[String]): Unit = {

    val spark = SparkSession.builder().appName("sparkToHive").enableHiveSupport().getOrCreate()
    import spark.implicits._

    val config = HBaseConfiguration.create()
    config.set("hbase.zookeeper.quorum", "ip's")
    config.set("hbase.zookeeper.property.clientPort","2181")
    config.set(TableInputFormat.INPUT_TABLE, "tableName")

    val newAPIJobConfiguration1 = Job.getInstance(config)
    newAPIJobConfiguration1.getConfiguration().set(TableOutputFormat.OUTPUT_TABLE, "tableName")
    newAPIJobConfiguration1.setOutputFormatClass(classOf[TableOutputFormat[ImmutableBytesWritable]])

    val df: DataFrame  = Seq(("foo", "1", "foo1"), ("bar", "2", "bar1")).toDF("key", "value1", "value2")

    val hbasePuts= df.rdd.map((row: Row) => {
      val  put = new Put(Bytes.toBytes(row.getString(0)))
      put.addColumn(Bytes.toBytes("cf1"), Bytes.toBytes("value1"), Bytes.toBytes(row.getString(1)))
      put.addColumn(Bytes.toBytes("cf2"), Bytes.toBytes("value2"), Bytes.toBytes(row.getString(2)))
      (new ImmutableBytesWritable(), put)
    })

    hbasePuts.saveAsNewAPIHadoopDataset(newAPIJobConfiguration1.getConfiguration())
    }

参考:https://sparkkb.wordpress.com/2015/05/04/save-javardd-to-hbase-using-saveasnewapihadoopdataset-spark-api-java-coding/

【讨论】:

  • 如果要保存到表中不应该是 config.set(TableInputFormat.OUPUT_TABLE, "tableName")
  • 这里使用的TableOutputFormat是一个HBase类文件。 hbase.apache.org/apidocs/org/apache/hadoop/hbase/mapreduce/… 因为我们想将数据推送到 HBase 表,我们正在设置 TableOutputFormat。 TableInputFormat 将具有 INPUT_TABLE 可以在我们从 HBase 中提取数据的情况下使用。
【解决方案2】:

下面是使用来自 Hortonworks 的 spark hbase 连接器的完整示例,该连接器位于 Maven

这个例子显示

  • 如何检查 HBase 表是否存在
  • 如果不存在则创建 HBase 表
  • 将 DataFrame 插入 HBase 表中
import org.apache.hadoop.hbase.client.{ColumnFamilyDescriptorBuilder, ConnectionFactory, TableDescriptorBuilder}
import org.apache.hadoop.hbase.{HBaseConfiguration, TableName}
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.execution.datasources.hbase.HBaseTableCatalog

object Main extends App {

  case class Employee(key: String, fName: String, lName: String, mName: String,
                      addressLine: String, city: String, state: String, zipCode: String)

  // as pre-requisites the table 'employee' with column families 'person' and 'address' should exist
  val tableNameString = "default:employee"
  val colFamilyPString = "person"
  val colFamilyAString = "address"
  val tableName = TableName.valueOf(tableNameString)
  val colFamilyP = colFamilyPString.getBytes
  val colFamilyA = colFamilyAString.getBytes

  val hBaseConf = HBaseConfiguration.create()
  val connection = ConnectionFactory.createConnection(hBaseConf);
  val admin = connection.getAdmin();

  println("Check if table 'employee' exists:")
  val tableExistsCheck: Boolean = admin.tableExists(tableName)
  println(s"Table " + tableName.toString + " exists? " + tableExistsCheck)

  if(tableExistsCheck == false) {
    println("Create Table employee with column families 'person' and 'address'")
    val colFamilyBuild1 = ColumnFamilyDescriptorBuilder.newBuilder(colFamilyP).build()
    val colFamilyBuild2 = ColumnFamilyDescriptorBuilder.newBuilder(colFamilyA).build()
    val tableDescriptorBuild = TableDescriptorBuilder.newBuilder(tableName)
      .setColumnFamily(colFamilyBuild1)
      .setColumnFamily(colFamilyBuild2)
      .build()
    admin.createTable(tableDescriptorBuild)
  }

  // define schema for the dataframe that should be loaded into HBase
  def catalog =
    s"""{
       |"table":{"namespace":"default","name":"employee"},
       |"rowkey":"key",
       |"columns":{
       |"key":{"cf":"rowkey","col":"key","type":"string"},
       |"fName":{"cf":"person","col":"firstName","type":"string"},
       |"lName":{"cf":"person","col":"lastName","type":"string"},
       |"mName":{"cf":"person","col":"middleName","type":"string"},
       |"addressLine":{"cf":"address","col":"addressLine","type":"string"},
       |"city":{"cf":"address","col":"city","type":"string"},
       |"state":{"cf":"address","col":"state","type":"string"},
       |"zipCode":{"cf":"address","col":"zipCode","type":"string"}
       |}
       |}""".stripMargin

  // define some test data
  val data = Seq(
    Employee("1","Horst","Hans","A","12main","NYC","NY","123"),
    Employee("2","Joe","Bill","B","1337ave","LA","CA","456"),
    Employee("3","Mohammed","Mohammed","C","1Apple","SanFran","CA","678")
  )

  // create SparkSession
  val spark: SparkSession = SparkSession.builder()
    .master("local[*]")
    .appName("HBaseConnector")
    .getOrCreate()

  // serialize data
  import spark.implicits._
  val df = spark.sparkContext.parallelize(data).toDF

  // write dataframe into HBase
  df.write.options(
    Map(HBaseTableCatalog.tableCatalog -> catalog, HBaseTableCatalog.newTable -> "3")) // create 3 regions
    .format("org.apache.spark.sql.execution.datasources.hbase")
    .save()

}

这对我有用,而我的 资源.

【讨论】:

    【解决方案3】:

    将答案用于代码格式化目的 医生告诉:

    sc.parallelize(data).toDF.write.options(
     Map(HBaseTableCatalog.tableCatalog -> catalog, HBaseTableCatalog.newTable -> "5"))
     .format("org.apache.hadoop.hbase.spark ")
     .save()
    

    sc.parallelize(data).toDF 是您的 DataFrame。文档示例使用 sc.parallelize(data).toDF

    将 scala 集合转换为数据帧

    你已经有了你的DataFrame,试着调用

    yourDataFrame.write.options(
         Map(HBaseTableCatalog.tableCatalog -> catalog, HBaseTableCatalog.newTable -> "5"))
         .format("org.apache.hadoop.hbase.spark ")
         .save()
    

    它应该可以工作。文档很清楚...

    UPD

    给定一个具有指定模式的 DataFrame,上面将创建一个 HBase 具有 5 个区域的表并将 DataFrame 保存在其中。请注意,如果 HBaseTableCatalog.newTable 没有指定,表必须是 预先创建的。

    这是关于数据分区的。每个 HBase 表可以有 1...X 个区域。您应该仔细选择区域数量。地区数量少是不好的。高地区数也不好。

    【讨论】:

    • 感谢您的回答;你能解释一下这一行吗:HBaseTableCatalog.newTable -&gt; "5"
    • 更新答案,见上文。 5,表示在HBase中为表创建5个区域
    • 目录在哪里定义? case class HBaseRecord( col0: String, col1: String, col2: String ) object HBaseRecord{ def apply(i: Int, t: String): HBaseRecord = { val s = s"""row${"%03d".format(i)}""" HBaseRecord(s, s"String$i: $t", s"String$i: $t") } } 之后怎么办?谢谢
    • 添加脚本后object HBaseRecord .我得到了这个错误error: too many arguments for method apply: (i: Int, t: String)HBaseRecord in object HBaseRecord &lt;console&gt;:1: error: ';' expected but 'for' found. 你能解释一下这个错误吗
    • 我不知道,最好提供完整的代码sn-p。错误非常明显。您以错误的方式调用 HBaseRecord.apply 记录。
    猜你喜欢
    • 2019-09-18
    • 1970-01-01
    • 1970-01-01
    • 2021-10-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-12-12
    • 1970-01-01
    相关资源
    最近更新 更多