-
- Downloads
[SPARK-18682][SS] Batch Source for Kafka
## What changes were proposed in this pull request? Today, you can start a stream that reads from kafka. However, given kafka's configurable retention period, it seems like sometimes you might just want to read all of the data that is available now. As such we should add a version that works with spark.read as well. The options should be the same as the streaming kafka source, with the following differences: startingOffsets should default to earliest, and should not allow latest (which would always be empty). endingOffsets should also be allowed and should default to latest. the same assign json format as startingOffsets should also be accepted. It would be really good, if things like .limit(n) were enough to prevent all the data from being read (this might just work). ## How was this patch tested? KafkaRelationSuite was added for testing batch queries via KafkaUtils. Author: Tyson Condie <tcondie@gmail.com> Closes #16686 from tcondie/SPARK-18682.
Showing
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/CachedKafkaConsumer.scala 73 additions, 29 deletions...a/org/apache/spark/sql/kafka010/CachedKafkaConsumer.scala
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/ConsumerStrategy.scala 84 additions, 0 deletions...cala/org/apache/spark/sql/kafka010/ConsumerStrategy.scala
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaOffsetRangeLimit.scala 51 additions, 0 deletions...org/apache/spark/sql/kafka010/KafkaOffsetRangeLimit.scala
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaOffsetReader.scala 312 additions, 0 deletions...ala/org/apache/spark/sql/kafka010/KafkaOffsetReader.scala
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaRelation.scala 124 additions, 0 deletions...n/scala/org/apache/spark/sql/kafka010/KafkaRelation.scala
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaSource.scala 41 additions, 282 deletions...ain/scala/org/apache/spark/sql/kafka010/KafkaSource.scala
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaSourceProvider.scala 183 additions, 79 deletions...a/org/apache/spark/sql/kafka010/KafkaSourceProvider.scala
- external/kafka-0-10-sql/src/main/scala/org/apache/spark/sql/kafka010/KafkaSourceRDD.scala 56 additions, 7 deletions.../scala/org/apache/spark/sql/kafka010/KafkaSourceRDD.scala
- external/kafka-0-10-sql/src/test/scala/org/apache/spark/sql/kafka010/KafkaRelationSuite.scala 233 additions, 0 deletions...la/org/apache/spark/sql/kafka010/KafkaRelationSuite.scala
- external/kafka-0-10-sql/src/test/scala/org/apache/spark/sql/kafka010/KafkaSourceSuite.scala 3 additions, 0 deletions...cala/org/apache/spark/sql/kafka010/KafkaSourceSuite.scala
- external/kafka-0-10-sql/src/test/scala/org/apache/spark/sql/kafka010/KafkaTestUtils.scala 20 additions, 1 deletion.../scala/org/apache/spark/sql/kafka010/KafkaTestUtils.scala
Loading
Please register or sign in to comment