-
- Downloads
Refactored streaming scheduler and added listener interface.
- Refactored Scheduler + JobManager to JobGenerator + JobScheduler and added JobSet for cleaner code. Moved scheduler related code to streaming.scheduler package. - Added StreamingListener trait (similar to SparkListener) to enable gathering to streaming stats like processing times and delays. StreamingContext.addListener() to added listeners. - Deduped some code in streaming tests by modifying TestSuiteBase, and added StreamingListenerSuite.
Showing
- core/src/main/scala/org/apache/spark/scheduler/SparkListener.scala 1 addition, 1 deletion...main/scala/org/apache/spark/scheduler/SparkListener.scala
- streaming/src/main/scala/org/apache/spark/streaming/Checkpoint.scala 1 addition, 1 deletion...rc/main/scala/org/apache/spark/streaming/Checkpoint.scala
- streaming/src/main/scala/org/apache/spark/streaming/DStream.scala 3 additions, 8 deletions...g/src/main/scala/org/apache/spark/streaming/DStream.scala
- streaming/src/main/scala/org/apache/spark/streaming/DStreamGraph.scala 1 addition, 0 deletions.../main/scala/org/apache/spark/streaming/DStreamGraph.scala
- streaming/src/main/scala/org/apache/spark/streaming/StreamingContext.scala 9 additions, 8 deletions...n/scala/org/apache/spark/streaming/StreamingContext.scala
- streaming/src/main/scala/org/apache/spark/streaming/dstream/ForEachDStream.scala 2 additions, 1 deletion...a/org/apache/spark/streaming/dstream/ForEachDStream.scala
- streaming/src/main/scala/org/apache/spark/streaming/dstream/NetworkInputDStream.scala 1 addition, 0 deletions.../apache/spark/streaming/dstream/NetworkInputDStream.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/BatchInfo.scala 38 additions, 0 deletions...cala/org/apache/spark/streaming/scheduler/BatchInfo.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/Job.scala 11 additions, 5 deletions...main/scala/org/apache/spark/streaming/scheduler/Job.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/JobGenerator.scala 22 additions, 26 deletions...a/org/apache/spark/streaming/scheduler/JobGenerator.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/JobScheduler.scala 104 additions, 0 deletions...a/org/apache/spark/streaming/scheduler/JobScheduler.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/JobSet.scala 61 additions, 0 deletions...n/scala/org/apache/spark/streaming/scheduler/JobSet.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/NetworkInputTracker.scala 2 additions, 1 deletion...pache/spark/streaming/scheduler/NetworkInputTracker.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/StreamingListener.scala 37 additions, 0 deletions.../apache/spark/streaming/scheduler/StreamingListener.scala
- streaming/src/main/scala/org/apache/spark/streaming/scheduler/StreamingListenerBus.scala 81 additions, 0 deletions...ache/spark/streaming/scheduler/StreamingListenerBus.scala
- streaming/src/test/scala/org/apache/spark/streaming/BasicOperationsSuite.scala 0 additions, 12 deletions...ala/org/apache/spark/streaming/BasicOperationsSuite.scala
- streaming/src/test/scala/org/apache/spark/streaming/CheckpointSuite.scala 10 additions, 16 deletions...st/scala/org/apache/spark/streaming/CheckpointSuite.scala
- streaming/src/test/scala/org/apache/spark/streaming/FailureSuite.scala 9 additions, 4 deletions.../test/scala/org/apache/spark/streaming/FailureSuite.scala
- streaming/src/test/scala/org/apache/spark/streaming/InputStreamsSuite.scala 0 additions, 12 deletions.../scala/org/apache/spark/streaming/InputStreamsSuite.scala
- streaming/src/test/scala/org/apache/spark/streaming/StreamingListenerSuite.scala 71 additions, 0 deletions...a/org/apache/spark/streaming/StreamingListenerSuite.scala
Loading
Please register or sign in to comment