-
- Downloads
[SPARK-4964][Streaming][Kafka] More updates to Exactly-once Kafka stream
Changes - Added example - Added a critical unit test that verifies that offset ranges can be recovered through checkpoints Might add more changes. Author: Tathagata Das <tathagata.das1565@gmail.com> Closes #4384 from tdas/new-kafka-fixes and squashes the following commits: 7c931c3 [Tathagata Das] Small update 3ed9284 [Tathagata Das] updated scala doc 83d0402 [Tathagata Das] Added JavaDirectKafkaWordCount example. 26df23c [Tathagata Das] Updates based on PR comments from Cody e4abf69 [Tathagata Das] Scala doc improvements and stuff. bb65232 [Tathagata Das] Fixed test bug and refactored KafkaStreamSuite 50f2b56 [Tathagata Das] Added Java API and added more Scala and Java unit tests. Also updated docs. e73589c [Tathagata Das] Minor changes. 4986784 [Tathagata Das] Added unit test to kafka offset recovery 6a91cab [Tathagata Das] Added example
Showing
- examples/scala-2.10/src/main/java/org/apache/spark/examples/streaming/JavaDirectKafkaWordCount.java 113 additions, 0 deletions...he/spark/examples/streaming/JavaDirectKafkaWordCount.java
- examples/scala-2.10/src/main/scala/org/apache/spark/examples/streaming/DirectKafkaWordCount.scala 71 additions, 0 deletions...pache/spark/examples/streaming/DirectKafkaWordCount.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/DirectKafkaInputDStream.scala 2 additions, 3 deletions...pache/spark/streaming/kafka/DirectKafkaInputDStream.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaCluster.scala 3 additions, 0 deletions...scala/org/apache/spark/streaming/kafka/KafkaCluster.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaRDD.scala 5 additions, 7 deletions...ain/scala/org/apache/spark/streaming/kafka/KafkaRDD.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaRDDPartition.scala 1 addition, 22 deletions.../org/apache/spark/streaming/kafka/KafkaRDDPartition.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaUtils.scala 274 additions, 79 deletions...n/scala/org/apache/spark/streaming/kafka/KafkaUtils.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/Leader.scala 16 additions, 5 deletions.../main/scala/org/apache/spark/streaming/kafka/Leader.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/OffsetRange.scala 47 additions, 6 deletions.../scala/org/apache/spark/streaming/kafka/OffsetRange.scala
- external/kafka/src/test/java/org/apache/spark/streaming/kafka/JavaDirectKafkaStreamSuite.java 159 additions, 0 deletions...che/spark/streaming/kafka/JavaDirectKafkaStreamSuite.java
- external/kafka/src/test/java/org/apache/spark/streaming/kafka/JavaKafkaStreamSuite.java 3 additions, 2 deletions...rg/apache/spark/streaming/kafka/JavaKafkaStreamSuite.java
- external/kafka/src/test/scala/org/apache/spark/streaming/kafka/DirectKafkaStreamSuite.scala 302 additions, 0 deletions...apache/spark/streaming/kafka/DirectKafkaStreamSuite.scala
- external/kafka/src/test/scala/org/apache/spark/streaming/kafka/KafkaClusterSuite.scala 10 additions, 14 deletions.../org/apache/spark/streaming/kafka/KafkaClusterSuite.scala
- external/kafka/src/test/scala/org/apache/spark/streaming/kafka/KafkaRDDSuite.scala 4 additions, 4 deletions...cala/org/apache/spark/streaming/kafka/KafkaRDDSuite.scala
- external/kafka/src/test/scala/org/apache/spark/streaming/kafka/KafkaStreamSuite.scala 36 additions, 26 deletions...a/org/apache/spark/streaming/kafka/KafkaStreamSuite.scala
- external/kafka/src/test/scala/org/apache/spark/streaming/kafka/ReliableKafkaStreamSuite.scala 2 additions, 2 deletions...ache/spark/streaming/kafka/ReliableKafkaStreamSuite.scala
Loading
Please register or sign in to comment