-
- Downloads
[SPARK-12177][STREAMING][KAFKA] Update KafkaDStreams to new Kafka 0.10 Consumer API
## What changes were proposed in this pull request? New Kafka consumer api for the released 0.10 version of Kafka ## How was this patch tested? Unit tests, manual tests Author: cody koeninger <cody@koeninger.org> Closes #11863 from koeninger/kafka-0.9.
Showing
- external/kafka-0-10-assembly/pom.xml 176 additions, 0 deletionsexternal/kafka-0-10-assembly/pom.xml
- external/kafka-0-10/pom.xml 98 additions, 0 deletionsexternal/kafka-0-10/pom.xml
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/CachedKafkaConsumer.scala 189 additions, 0 deletions...apache/spark/streaming/kafka010/CachedKafkaConsumer.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/ConsumerStrategy.scala 314 additions, 0 deletions...rg/apache/spark/streaming/kafka010/ConsumerStrategy.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/DirectKafkaInputDStream.scala 318 additions, 0 deletions...he/spark/streaming/kafka010/DirectKafkaInputDStream.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/KafkaRDD.scala 232 additions, 0 deletions.../scala/org/apache/spark/streaming/kafka010/KafkaRDD.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/KafkaRDDPartition.scala 45 additions, 0 deletions...g/apache/spark/streaming/kafka010/KafkaRDDPartition.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/KafkaTestUtils.scala 277 additions, 0 deletions.../org/apache/spark/streaming/kafka010/KafkaTestUtils.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/KafkaUtils.scala 175 additions, 0 deletions...cala/org/apache/spark/streaming/kafka010/KafkaUtils.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/LocationStrategy.scala 77 additions, 0 deletions...rg/apache/spark/streaming/kafka010/LocationStrategy.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/OffsetRange.scala 153 additions, 0 deletions...ala/org/apache/spark/streaming/kafka010/OffsetRange.scala
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/package-info.java 21 additions, 0 deletions...ala/org/apache/spark/streaming/kafka010/package-info.java
- external/kafka-0-10/src/main/scala/org/apache/spark/streaming/kafka010/package.scala 23 additions, 0 deletions...n/scala/org/apache/spark/streaming/kafka010/package.scala
- external/kafka-0-10/src/test/java/org/apache/spark/streaming/kafka010/JavaConsumerStrategySuite.java 84 additions, 0 deletions...e/spark/streaming/kafka010/JavaConsumerStrategySuite.java
- external/kafka-0-10/src/test/java/org/apache/spark/streaming/kafka010/JavaDirectKafkaStreamSuite.java 180 additions, 0 deletions.../spark/streaming/kafka010/JavaDirectKafkaStreamSuite.java
- external/kafka-0-10/src/test/java/org/apache/spark/streaming/kafka010/JavaKafkaRDDSuite.java 122 additions, 0 deletions...rg/apache/spark/streaming/kafka010/JavaKafkaRDDSuite.java
- external/kafka-0-10/src/test/java/org/apache/spark/streaming/kafka010/JavaLocationStrategySuite.java 58 additions, 0 deletions...e/spark/streaming/kafka010/JavaLocationStrategySuite.java
- external/kafka-0-10/src/test/resources/log4j.properties 28 additions, 0 deletionsexternal/kafka-0-10/src/test/resources/log4j.properties
- external/kafka-0-10/src/test/scala/org/apache/spark/streaming/kafka010/DirectKafkaStreamSuite.scala 612 additions, 0 deletions...che/spark/streaming/kafka010/DirectKafkaStreamSuite.scala
- external/kafka-0-10/src/test/scala/org/apache/spark/streaming/kafka010/KafkaRDDSuite.scala 169 additions, 0 deletions...a/org/apache/spark/streaming/kafka010/KafkaRDDSuite.scala
Loading
Please register or sign in to comment