Skip to content

Commit d904b25

Browse files
committed
Remove duplicate
1 parent ec2b6c2 commit d904b25

File tree

1 file changed

+2
-1
lines changed

1 file changed

+2
-1
lines changed

streaming/src/main/scala/org/apache/spark/streaming/scheduler/JobGenerator.scala

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -220,7 +220,8 @@ class JobGenerator(jobScheduler: JobScheduler) extends Logging {
220220
logInfo("Batches pending processing (" + pendingTimes.size + " batches): " +
221221
pendingTimes.mkString(", "))
222222
// Reschedule jobs for these times
223-
val timesToReschedule = (pendingTimes ++ downTimes).distinct.sorted(Time.ordering)
223+
val timesToReschedule = (pendingTimes ++ downTimes).filter { _ != restartTime }
224+
.distinct.sorted(Time.ordering)
224225
logInfo("Batches to reschedule (" + timesToReschedule.size + " batches): " +
225226
timesToReschedule.mkString(", "))
226227
timesToReschedule.foreach { time =>

0 commit comments

Comments
 (0)