From 10ea8e761bec796438fb75b2df8933181385a318 Mon Sep 17 00:00:00 2001 From: Scott Gorlin Date: Sat, 16 Jan 2016 22:09:22 -0500 Subject: [PATCH 1/7] Upgrade deps to hbase 1.0.0 --- build.sbt | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/build.sbt b/build.sbt index 0c68167..6f1ce11 100644 --- a/build.sbt +++ b/build.sbt @@ -1,16 +1,16 @@ lazy val root = (project in file(".")). settings( name := "spark_hbase", - version := "1.0", + version := "1.1", scalaVersion := "2.10.5", sparkVersion := "1.3.0" ) libraryDependencies ++= Seq( - "org.apache.hbase" % "hbase" % "0.98.8-hadoop2", - "org.apache.hbase" % "hbase-client" % "0.98.8-hadoop2", - "org.apache.hbase" % "hbase-common" % "0.98.8-hadoop2", - "org.apache.hbase" % "hbase-server" % "0.98.8-hadoop2" + "org.apache.hbase" % "hbase" % "1.0.0", + "org.apache.hbase" % "hbase-client" % "1.0.0", + "org.apache.hbase" % "hbase-common" % "1.0.0", + "org.apache.hbase" % "hbase-server" % "1.0.0" ) mergeStrategy in assembly := { From fff36f0b73a0ac972fcbd393bc60f8cf6320fe9a Mon Sep 17 00:00:00 2001 From: Scott Gorlin Date: Sat, 16 Jan 2016 22:10:58 -0500 Subject: [PATCH 2/7] HBaseResultToMapConverter: latest columns/values direct to python dict --- .../scala/examples/pythonConverters.scala | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/src/main/scala/examples/pythonConverters.scala b/src/main/scala/examples/pythonConverters.scala index 9a3cf0c..22b61a8 100644 --- a/src/main/scala/examples/pythonConverters.scala +++ b/src/main/scala/examples/pythonConverters.scala @@ -54,6 +54,27 @@ class HBaseResultToStringConverter extends Converter[Any, String]{ } } +/* Map of "columnFamily:column"->"value" consistent with Python dict + Only works with 1 version max. Ser/deser as python dict naturally, via HashMap + */ +class HBaseResultToMapConverter extends Converter[Any, java.util.Map[String, String]]{ + override def convert(obj: Any): java.util.Map[String, String] = { + + import collection.JavaConverters._ + val result = obj.asInstanceOf[Result] + + val my_map = result.listCells.asScala.map( + cell => ( + Bytes.toString(CellUtil.cloneFamily(cell)) + ":" + + Bytes.toString(CellUtil.cloneQualifier(cell)), + Bytes.toString(CellUtil.cloneValue(cell))) + ).toMap.asJava + + // nb: java.util.Map and scala map do not serialize into python?? + new java.util.HashMap[String, String](my_map) + } +} + /** * Implementation of [[org.apache.spark.api.python.Converter]] that converts an * ImmutableBytesWritable to a String From f21cbb3647e5b4a4988000689d3d044bacce6f12 Mon Sep 17 00:00:00 2001 From: Scott Gorlin Date: Sat, 16 Jan 2016 22:12:03 -0500 Subject: [PATCH 3/7] HBaseRawResultsConverter: list of cells direct to python --- src/main/scala/examples/pythonConverters.scala | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/src/main/scala/examples/pythonConverters.scala b/src/main/scala/examples/pythonConverters.scala index 22b61a8..07a211a 100644 --- a/src/main/scala/examples/pythonConverters.scala +++ b/src/main/scala/examples/pythonConverters.scala @@ -54,6 +54,24 @@ class HBaseResultToStringConverter extends Converter[Any, String]{ } } +class HBaseRawResultsConverter extends Converter[Any, java.util.List[java.util.Map[String, String]]]{ + override def convert(obj: Any): java.util.List[java.util.Map[String, String]] = { + import collection.JavaConverters._ + val result = obj.asInstanceOf[Result] + val output = result.listCells.asScala.map(cell => + new java.util.HashMap[String, String](Map( + "row" -> Bytes.toString(CellUtil.cloneFamily(cell)), + "columnFamily" -> Bytes.toString(CellUtil.cloneFamily(cell)), + "qualifier" -> Bytes.toString(CellUtil.cloneQualifier(cell)), + "timestamp" -> cell.getTimestamp.toString, + "type" -> Type.codeToType(cell.getTypeByte).toString, + "value" -> Bytes.toString(CellUtil.cloneValue(cell)) + ).asJava) + ) + new java.util.ArrayList[java.util.Map[String, String]](output.asJava) + } +} + /* Map of "columnFamily:column"->"value" consistent with Python dict Only works with 1 version max. Ser/deser as python dict naturally, via HashMap */ From d060fb63fc53224129909f1162b935e77cd7d343 Mon Sep 17 00:00:00 2001 From: Scott Gorlin Date: Sat, 16 Jan 2016 22:13:25 -0500 Subject: [PATCH 4/7] MaxHBaseTimestamp: converter which only returns the last timestamp per result --- src/main/scala/examples/pythonConverters.scala | 11 +++++++++++ 1 file changed, 11 insertions(+) diff --git a/src/main/scala/examples/pythonConverters.scala b/src/main/scala/examples/pythonConverters.scala index 07a211a..da5d433 100644 --- a/src/main/scala/examples/pythonConverters.scala +++ b/src/main/scala/examples/pythonConverters.scala @@ -93,6 +93,17 @@ class HBaseResultToMapConverter extends Converter[Any, java.util.Map[String, Str } } +/* Returns the most recent cell timestamp on the row */ +class MaxHBaseTimestamp extends Converter[Any, Long]{ + override def convert(obj: Any): Long = { + import collection.JavaConverters._ + val result = obj.asInstanceOf[Result] + result.listCells.asScala.map(cell => + cell.getTimestamp + ).max + } +} + /** * Implementation of [[org.apache.spark.api.python.Converter]] that converts an * ImmutableBytesWritable to a String From fc42528d1317776a0d987d46b2cac11a90ed7799 Mon Sep 17 00:00:00 2001 From: Scott Gorlin Date: Sat, 16 Jan 2016 22:19:18 -0500 Subject: [PATCH 5/7] HBaseResultToJSONConverter: returns valid JSON array --- .../scala/examples/pythonConverters.scala | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/src/main/scala/examples/pythonConverters.scala b/src/main/scala/examples/pythonConverters.scala index da5d433..6b9d479 100644 --- a/src/main/scala/examples/pythonConverters.scala +++ b/src/main/scala/examples/pythonConverters.scala @@ -54,6 +54,25 @@ class HBaseResultToStringConverter extends Converter[Any, String]{ } } +class HBaseResultToJSONConverter extends Converter[Any, String]{ + override def convert(obj: Any): String = { + import collection.JavaConverters._ + val result = obj.asInstanceOf[Result] + val output = result.listCells.asScala.map(cell => + Map( + "row" -> Bytes.toStringBinary(CellUtil.cloneFamily(cell)), + "columnFamily" -> Bytes.toStringBinary(CellUtil.cloneFamily(cell)), + "qualifier" -> Bytes.toStringBinary(CellUtil.cloneQualifier(cell)), + "timestamp" -> cell.getTimestamp.toString, + "type" -> Type.codeToType(cell.getTypeByte).toString, + "value" -> Bytes.toStringBinary(CellUtil.cloneValue(cell)) + ) + ) + // output is a JSON array of objects (maps) + "[" + output.map(JSONObject(_).toString()).mkString(", ") + "]" + } +} + class HBaseRawResultsConverter extends Converter[Any, java.util.List[java.util.Map[String, String]]]{ override def convert(obj: Any): java.util.List[java.util.Map[String, String]] = { import collection.JavaConverters._ From 5685555cb01f9c7879e55394ed16544c8a6c0c5d Mon Sep 17 00:00:00 2001 From: Scott Gorlin Date: Mon, 25 Jan 2016 16:54:08 -0500 Subject: [PATCH 6/7] DOC: conform to existing doc standard and add a few lines --- .../scala/examples/pythonConverters.scala | 23 +++++++++++++++---- 1 file changed, 18 insertions(+), 5 deletions(-) diff --git a/src/main/scala/examples/pythonConverters.scala b/src/main/scala/examples/pythonConverters.scala index 6b9d479..0a4c638 100644 --- a/src/main/scala/examples/pythonConverters.scala +++ b/src/main/scala/examples/pythonConverters.scala @@ -31,7 +31,7 @@ import org.apache.hadoop.hbase.CellUtil /** * Implementation of [[org.apache.spark.api.python.Converter]] that converts all * the records in an HBase Result to a String. In the String, it contains row, column, - * qualifier, timesstamp, type and value + * qualifier, timestamp, type and value */ class HBaseResultToStringConverter extends Converter[Any, String]{ @@ -54,6 +54,11 @@ class HBaseResultToStringConverter extends Converter[Any, String]{ } } +/** + * Implementation of [[org.apache.spark.api.python.Converter]] that converts all + * the records in an HBase Result to a String. In the String, it contains row, column, + * qualifier, timestamp, type and value, all packed as JSON + */ class HBaseResultToJSONConverter extends Converter[Any, String]{ override def convert(obj: Any): String = { import collection.JavaConverters._ @@ -73,6 +78,11 @@ class HBaseResultToJSONConverter extends Converter[Any, String]{ } } +/** + * Implementation of [[org.apache.spark.api.python.Converter]] that converts all + * the records in an HBase Result into a list of maps, each containing row, column, + * qualifier, timestamp, type and value + */ class HBaseRawResultsConverter extends Converter[Any, java.util.List[java.util.Map[String, String]]]{ override def convert(obj: Any): java.util.List[java.util.Map[String, String]] = { import collection.JavaConverters._ @@ -91,8 +101,9 @@ class HBaseRawResultsConverter extends Converter[Any, java.util.List[java.util.M } } -/* Map of "columnFamily:column"->"value" consistent with Python dict - Only works with 1 version max. Ser/deser as python dict naturally, via HashMap +/** + * Map of "columnFamily:column"->"value" consistent with Python dict + * Only works with 1 version max. Ser/deser as python dict naturally, via HashMap */ class HBaseResultToMapConverter extends Converter[Any, java.util.Map[String, String]]{ override def convert(obj: Any): java.util.Map[String, String] = { @@ -112,7 +123,10 @@ class HBaseResultToMapConverter extends Converter[Any, java.util.Map[String, Str } } -/* Returns the most recent cell timestamp on the row */ +/** + * Returns the most recent cell timestamp on the row as a long integer, useful for + * operations that require syncing to data changes. + */ class MaxHBaseTimestamp extends Converter[Any, Long]{ override def convert(obj: Any): Long = { import collection.JavaConverters._ @@ -127,7 +141,6 @@ class MaxHBaseTimestamp extends Converter[Any, Long]{ * Implementation of [[org.apache.spark.api.python.Converter]] that converts an * ImmutableBytesWritable to a String */ - class ImmutableBytesWritableToStringConverter extends Converter[Any, String] { override def convert(obj: Any): String = { val key = obj.asInstanceOf[ImmutableBytesWritable] From df45276a744e856e18922fecdb2c626664918d5f Mon Sep 17 00:00:00 2001 From: Scott Gorlin Date: Mon, 25 Jan 2016 17:04:54 -0500 Subject: [PATCH 7/7] DOC and a few renames for clarity --- src/main/scala/examples/pythonConverters.scala | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/src/main/scala/examples/pythonConverters.scala b/src/main/scala/examples/pythonConverters.scala index 0a4c638..31df49e 100644 --- a/src/main/scala/examples/pythonConverters.scala +++ b/src/main/scala/examples/pythonConverters.scala @@ -81,9 +81,10 @@ class HBaseResultToJSONConverter extends Converter[Any, String]{ /** * Implementation of [[org.apache.spark.api.python.Converter]] that converts all * the records in an HBase Result into a list of maps, each containing row, column, - * qualifier, timestamp, type and value + * qualifier, timestamp, type and value. Values are returned as raw bytes rather than as + * printable versions, and thus may require unpacking. */ -class HBaseRawResultsConverter extends Converter[Any, java.util.List[java.util.Map[String, String]]]{ +class HBaseResultToListConverter extends Converter[Any, java.util.List[java.util.Map[String, String]]]{ override def convert(obj: Any): java.util.List[java.util.Map[String, String]] = { import collection.JavaConverters._ val result = obj.asInstanceOf[Result] @@ -102,8 +103,15 @@ class HBaseRawResultsConverter extends Converter[Any, java.util.List[java.util.M } /** - * Map of "columnFamily:column"->"value" consistent with Python dict - * Only works with 1 version max. Ser/deser as python dict naturally, via HashMap + * Implementation of [[org.apache.spark.api.python.Converter]] that converts all + * the records in an HBase Result into a Map of "columnFamily:column"->"value" consistent + * with a Python dict. + * + * Only works with 1 version max (possibly latest if multiple are returned?). + * Ser/deser as python dict naturally, via HashMap. + * + * Values are returned as raw bytes rather than as printable versions, and thus + * may require unpacking. */ class HBaseResultToMapConverter extends Converter[Any, java.util.Map[String, String]]{ override def convert(obj: Any): java.util.Map[String, String] = {