Home
last modified time | relevance | path

Searched refs:dagScheduler (Results 1 – 13 of 13) sorted by relevance

/dports/devel/spark/spark-2.1.1/core/src/main/scala/org/apache/spark/scheduler/
H A DDAGSchedulerSource.scala24 private[scheduler] class DAGSchedulerSource(val dagScheduler: DAGScheduler)
30 override def getValue: Int = dagScheduler.failedStages.size
34 override def getValue: Int = dagScheduler.runningStages.size
38 override def getValue: Int = dagScheduler.waitingStages.size
42 override def getValue: Int = dagScheduler.numTotalJobs
46 override def getValue: Int = dagScheduler.activeJobs.size
H A DDAGScheduler.scala1619 dagScheduler.handleStageCancellation(stageId)
1622 dagScheduler.handleJobCancellation(jobId)
1625 dagScheduler.handleJobGroupCancelled(groupId)
1628 dagScheduler.doCancelAllJobs()
1641 dagScheduler.handleBeginEvent(task, taskInfo)
1644 dagScheduler.handleGetTaskResult(taskInfo)
1647 dagScheduler.handleTaskCompletion(completion)
1653 dagScheduler.resubmitFailedStages()
1659 dagScheduler.doCancelAllJobs()
1663 dagScheduler.sc.stopInNewThread()
[all …]
H A DTaskSchedulerImpl.scala112 var dagScheduler: DAGScheduler = null variable
132 override def setDAGScheduler(dagScheduler: DAGScheduler) {
133 this.dagScheduler = dagScheduler
391 dagScheduler.executorLost(failedExecutor.get, reason.get)
414 dagScheduler.executorHeartbeatReceived(execId, accumUpdatesWithTaskIds, blockManagerId)
517 dagScheduler.executorLost(failedExecutor.get, reason)
580 dagScheduler.executorAdded(execId, host)
H A DJobWaiter.scala31 dagScheduler: DAGScheduler,
53 dagScheduler.cancelJob(jobId)
H A DTaskSetManager.scala478 sched.dagScheduler.taskStarted(task, info)
647 sched.dagScheduler.taskGettingResult(info)
682 sched.dagScheduler.taskEnded(tasks(index), Success, result.value(), result.accumUpdates, info)
782 sched.dagScheduler.taskEnded(tasks(index), reason, null, accumUpdates, info)
811 sched.dagScheduler.taskSetFailed(taskSet, message, exception)
863 sched.dagScheduler.taskEnded(
H A DTaskScheduler.scala58 def setDAGScheduler(dagScheduler: DAGScheduler): Unit
/dports/devel/spark/spark-2.1.1/core/src/test/scala/org/apache/spark/scheduler/
H A DTaskSchedulerImplSuite.scala47 var dagScheduler: DAGScheduler = null variable
62 if (dagScheduler != null) {
63 dagScheduler.stop()
64 dagScheduler = null
77 dagScheduler = new DAGScheduler(sc, taskScheduler) {
H A DExternalClusterManagerSuite.scala82 override def setDAGScheduler(dagScheduler: DAGScheduler): Unit = {}
H A DOutputCommitCoordinatorSuite.scala127 sc.dagScheduler = dagSchedulerWithMockTaskScheduler
H A DTaskResultGetterSuite.scala227 sc.dagScheduler = newDAGScheduler
H A DDAGSchedulerSuite.scala39 class DAGSchedulerEventProcessLoopTester(dagScheduler: DAGScheduler)
40 extends DAGSchedulerEventProcessLoop(dagScheduler) {
53 dagScheduler.stop()
129 override def setDAGScheduler(dagScheduler: DAGScheduler) = {}
555 override def setDAGScheduler(dagScheduler: DAGScheduler): Unit = {}
H A DTaskSetManagerSuite.scala100 dagScheduler = new FakeDAGScheduler(sc, this)
172 sched.dagScheduler.stop()
/dports/devel/spark/spark-2.1.1/core/src/main/scala/org/apache/spark/
H A DSparkContext.scala304 private[spark] def dagScheduler: DAGScheduler = _dagScheduler
1690 dagScheduler.getPreferredLocs(rdd, partition)
1925 dagScheduler.runJob(rdd, cleanedFunc, partitions, callSite, resultHandler, localProperties.get)
2006 val result = dagScheduler.runApproximateJob(rdd, cleanedFunc, evaluator, callSite, timeout,
2026 val waiter = dagScheduler.submitJob(
2045 val waiter = dagScheduler.submitMapStage(
2059 dagScheduler.cancelJobGroup(groupId)
2065 dagScheduler.cancelAllJobs()
2075 dagScheduler.cancelJob(jobId)
2085 dagScheduler.cancelStage(stageId)