-
- Downloads
[SPARK-8127] [STREAMING] [KAFKA] KafkaRDD optimize count() take() isEmpty()
…ed KafkaRDD methods. Possible fix for [SPARK-7122], but probably a worthwhile optimization regardless. Author: cody koeninger <cody@koeninger.org> Closes #6632 from koeninger/kafka-rdd-count and squashes the following commits: 321340d [cody koeninger] [SPARK-8127][Streaming][Kafka] additional test of ordering of take() 5a05d0f [cody koeninger] [SPARK-8127][Streaming][Kafka] additional test of isEmpty f68bd32 [cody koeninger] [Streaming][Kafka][SPARK-8127] code cleanup 9555b73 [cody koeninger] Merge branch 'master' into kafka-rdd-count 253031d [cody koeninger] [Streaming][Kafka][SPARK-8127] mima exclusion for change to private method 8974b9e [cody koeninger] [Streaming][Kafka][SPARK-8127] check offset ranges before constructing KafkaRDD c3768c5 [cody koeninger] [Streaming][Kafka] Take advantage of offset range info for size-related KafkaRDD methods. Possible fix for [SPARK-7122], but probably a worthwhile optimization regardless.
Showing
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/DirectKafkaInputDStream.scala 2 additions, 6 deletions...pache/spark/streaming/kafka/DirectKafkaInputDStream.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaCluster.scala 8 additions, 0 deletions...scala/org/apache/spark/streaming/kafka/KafkaCluster.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaRDD.scala 44 additions, 0 deletions...ain/scala/org/apache/spark/streaming/kafka/KafkaRDD.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaRDDPartition.scala 4 additions, 1 deletion.../org/apache/spark/streaming/kafka/KafkaRDDPartition.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/KafkaUtils.scala 32 additions, 14 deletions...n/scala/org/apache/spark/streaming/kafka/KafkaUtils.scala
- external/kafka/src/main/scala/org/apache/spark/streaming/kafka/OffsetRange.scala 6 additions, 0 deletions.../scala/org/apache/spark/streaming/kafka/OffsetRange.scala
- external/kafka/src/test/scala/org/apache/spark/streaming/kafka/KafkaRDDSuite.scala 23 additions, 3 deletions...cala/org/apache/spark/streaming/kafka/KafkaRDDSuite.scala
- project/MimaExcludes.scala 3 additions, 0 deletionsproject/MimaExcludes.scala
Loading
Please register or sign in to comment