【问题标题】:Creating an RDD[(ImmutableBytesWritable, Result)] in Scala在 Scala 中创建 RDD[(ImmutableBytesWritable, Result)]
【发布时间】:2018-08-07 18:45:00
【问题描述】:

我正在做一个单元测试,我需要创建一个RDD[(ImmutableBytesWritable, Result)]。数据只包含一个唯一的id 和一个非唯一的value 列。

我可以使用toDF 从列表中创建一个DataFrame,并将其转换为RDD[Row];但我无法将其映射到 RDD[(ImmutableBytesWritable, Result)]

val values = List((1, 1234), (2, 123), (3, 1234))
import spark.implicits._
val df = values.toDF("id", "value")
val counts : RDD[(ImmutableBytesWritable, Result)] = df.rdd.map(
  row => (new ImmutableBytesWritable(), Result.create(...))
)

谢谢!

【问题讨论】:

  • 你能发布你到目前为止的代码吗?
  • @Jeremy 是的,好电话。
  • 什么是结果?
  • @Jeremy 这是org.apache.hadoop.hbase.client.Result

标签: scala apache-spark hbase rdd


【解决方案1】:

我不使用 HBase,但考虑到您正在尝试构建的签名,我想出了这个。试试看吧。

我使用 CellUtil 创建了一个 Cell 来构建 Result。

import org.apache.hadoop.hbase.{Cell, CellUtil}

import scala.collection.JavaConversions._
import scala.collection.mutable.ListBuffer

import scala.math.BigInt

import org.apache.spark._
import org.apache.spark.rdd._
import org.apache.spark.sql._

import org.apache.hadoop.hbase.client.Result
import org.apache.hadoop.hbase.io.ImmutableBytesWritable

object Practice extends App {

  val sparkConfig = new SparkConf().setAppName("test").setMaster("local[*]")
  val ss = SparkSession.builder().config(sparkConfig).getOrCreate()
  val values = Seq((1, 1234), (2, 123), (3, 1234))

  import ss.implicits._

  val df = values.toDF("id", "value")
  val counts: RDD[(ImmutableBytesWritable, Result)] = df.rdd.map{ row =>
    val key = row.getAs[Int]("id")
    val keyByteArray = BigInt(key).toByteArray
    val ibw = new ImmutableBytesWritable()
    ibw.set(keyByteArray)

    val value = row.getAs[Int]("value")
    val valueByteArray = BigInt(value).toByteArray
    val cellList = List(CellUtil.createCell(valueByteArray))
    val cell: java.util.List[Cell] = ListBuffer(cellList: _*)
    val result = Result.create(cell)

    (ibw, result)

  }
}

打印结果,这并不代表一个好的答案,给你这个:

KeyValue(03,keyvalues={\x04\xD2//LATEST_TIMESTAMP/Maximum/vlen=0/seqid=0})
KeyValue(01,keyvalues={\x04\xD2//LATEST_TIMESTAMP/Maximum/vlen=0/seqid=0})
KeyValue(02,keyvalues={{//LATEST_TIMESTAMP/Maximum/vlen=0/seqid=0})

【讨论】:

  • 请注意,我从 Scala List 转换为 java.util.List。您可能可以跳过它,只需创建一个 Java 列表即可。
  • 谢谢!你想对打印的结果说什么?
  • 没什么特别的。只是显示应用程序运行。由于我没有 HBase,我实际上无法测试结果。如果您发现我的回答有帮助,请考虑将我的回答标记为“已接受”。
猜你喜欢
  • 2019-11-29
  • 1970-01-01
  • 2019-08-07
  • 2018-03-24
  • 1970-01-01
  • 2019-08-11
  • 2016-05-06
  • 2018-06-16
  • 1970-01-01
相关资源
最近更新 更多