[SPARK-20217][CORE] Executor should not fail stage if killed task throws non-interrupted exception
## What changes were proposed in this pull request? If tasks throw non-interrupted exceptions on kill (e.g. java.nio.channels.ClosedByInterruptException), their death is reported back as TaskFailed instead of TaskKilled. This causes stage failure in some cases. This is reproducible as follows. Run the following, and then use SparkContext.killTaskAttempt to kill one of the tasks. The entire stage will fail since we threw a RuntimeException instead of InterruptedException. ``` spark.range(100).repartition(100).foreach { i => try { Thread.sleep(10000000) } catch { case t: InterruptedException => throw new RuntimeException(t) } } ``` Based on the code in TaskSetManager, I think this also affects kills of speculative tasks. However, since the number of speculated tasks is few, and usually you need to fail a task a few times before the stage is cancelled, it unlikely this would be noticed in production unless both speculation was enabled and the num allowed task failures was = 1. We should probably unconditionally return TaskKilled instead of TaskFailed if the task was killed by the driver, regardless of the actual exception thrown. ## How was this patch tested? Unit test. The test fails before the change in Executor.scala cc JoshRosen Author: Eric Liang <ekl@databricks.com> Closes #17531 from ericl/fix-task-interrupt.
This commit is contained in:
parent
4000f128b7
commit
5142e5d4e0
|
@ -432,7 +432,7 @@ private[spark] class Executor(
|
|||
setTaskFinishedAndClearInterruptStatus()
|
||||
execBackend.statusUpdate(taskId, TaskState.KILLED, ser.serialize(TaskKilled(t.reason)))
|
||||
|
||||
case _: InterruptedException if task.reasonIfKilled.isDefined =>
|
||||
case NonFatal(_) if task != null && task.reasonIfKilled.isDefined =>
|
||||
val killReason = task.reasonIfKilled.getOrElse("unknown reason")
|
||||
logInfo(s"Executor interrupted and killed $taskName (TID $taskId), reason: $killReason")
|
||||
setTaskFinishedAndClearInterruptStatus()
|
||||
|
|
|
@ -572,7 +572,13 @@ class SparkContextSuite extends SparkFunSuite with LocalSparkContext with Eventu
|
|||
// first attempt will hang
|
||||
if (!SparkContextSuite.isTaskStarted) {
|
||||
SparkContextSuite.isTaskStarted = true
|
||||
Thread.sleep(9999999)
|
||||
try {
|
||||
Thread.sleep(9999999)
|
||||
} catch {
|
||||
case t: Throwable =>
|
||||
// SPARK-20217 should not fail stage if task throws non-interrupted exception
|
||||
throw new RuntimeException("killed")
|
||||
}
|
||||
}
|
||||
// second attempt succeeds immediately
|
||||
}
|
||||
|
|
Loading…
Reference in a new issue