Hao Hao has posted comments on this change. ( http://gerrit.cloudera.org:8080/8552 )
Change subject: KUDU-1454 [part 1]: update propagated timestamp on the driver ...................................................................... Patch Set 3: (8 comments) http://gerrit.cloudera.org:8080/#/c/8552/1//COMMIT_MSG Commit Message: http://gerrit.cloudera.org:8080/#/c/8552/1//COMMIT_MSG@9 PS1, Line 9: Currently, Spark uses multiple clients simultaneously, possibly on > typo Done http://gerrit.cloudera.org:8080/#/c/8552/1//COMMIT_MSG@16 PS1, Line 16: This patch uses Accumulator in Spark to properly update the propagated > typo: patch Done http://gerrit.cloudera.org:8080/#/c/8552/2//COMMIT_MSG Commit Message: http://gerrit.cloudera.org:8080/#/c/8552/2//COMMIT_MSG@9 PS2, Line 9: Spark > Spark Done http://gerrit.cloudera.org:8080/#/c/8552/1/java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala File java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala: http://gerrit.cloudera.org:8080/#/c/8552/1/java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala@58 PS1, Line 58: private[kudu] class TimestampAccumulator(var timestamp: Long = 0L) > should this be marked private somehow? (I'm not good at Scala) Done http://gerrit.cloudera.org:8080/#/c/8552/1/java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala@66 PS1, Line 66: } : : override def reset(): Unit = { : time > eg here you can just use = new TimestampAccumulator(timestamp) Done http://gerrit.cloudera.org:8080/#/c/8552/1/java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala@294 PS1, Line 294: case DataTypes.TimestampType => operation.getRow.addLong(kuduIdx, KuduRelation.timestampToMicros(row.getTimestamp(sparkIdx))) > I do not quite follow your second comment. My understanding is for each RDD After discussed with Dan, figured out Todd's point is valid and need to propagate back the last propagated timestamp from the driver to executors. Updated accordingly. http://gerrit.cloudera.org:8080/#/c/8552/2/java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala File java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala: http://gerrit.cloudera.org:8080/#/c/8552/2/java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala@33 PS2, Line 33: import org.apache.spark.util.AccumulatorV2 > The braces are unnecessary with just one class Done http://gerrit.cloudera.org:8080/#/c/8552/2/java/kudu-spark/src/main/scala/org/apache/kudu/spark/kudu/KuduContext.scala@58 PS2, Line 58: private[kudu] class TimestampAccumulator(var timestamp: Long = 0L) > make private Done -- To view, visit http://gerrit.cloudera.org:8080/8552 To unsubscribe, visit http://gerrit.cloudera.org:8080/settings Gerrit-Project: kudu Gerrit-Branch: master Gerrit-MessageType: comment Gerrit-Change-Id: Id0a078ae8ebaa6a859be75c822879291192c5842 Gerrit-Change-Number: 8552 Gerrit-PatchSet: 3 Gerrit-Owner: Hao Hao <[email protected]> Gerrit-Reviewer: Dan Burkert <[email protected]> Gerrit-Reviewer: Hao Hao <[email protected]> Gerrit-Reviewer: Kudu Jenkins Gerrit-Reviewer: Todd Lipcon <[email protected]> Gerrit-Comment-Date: Thu, 16 Nov 2017 23:46:19 +0000 Gerrit-HasComments: Yes
