public class ExecutionVertex extends Object implements Serializable
Execution.| Constructor and Description |
|---|
ExecutionVertex(ExecutionJobVertex jobVertex,
int subTaskIndex,
IntermediateResult[] producedDataSets,
scala.concurrent.duration.FiniteDuration timeout) |
ExecutionVertex(ExecutionJobVertex jobVertex,
int subTaskIndex,
IntermediateResult[] producedDataSets,
scala.concurrent.duration.FiniteDuration timeout,
long createTimestamp) |
public ExecutionVertex(ExecutionJobVertex jobVertex, int subTaskIndex, IntermediateResult[] producedDataSets, scala.concurrent.duration.FiniteDuration timeout)
public ExecutionVertex(ExecutionJobVertex jobVertex, int subTaskIndex, IntermediateResult[] producedDataSets, scala.concurrent.duration.FiniteDuration timeout, long createTimestamp)
public JobID getJobId()
public ExecutionJobVertex getJobVertex()
public JobVertexID getJobvertexId()
public String getTaskName()
public String getTaskNameWithSubtaskIndex()
public int getTotalNumberOfParallelSubtasks()
public int getParallelSubtaskIndex()
public int getNumberOfInputs()
public ExecutionEdge[] getInputEdges(int input)
public CoLocationConstraint getLocationConstraint()
public Execution getCurrentExecutionAttempt()
public ExecutionState getExecutionState()
public long getStateTimestamp(ExecutionState state)
public Throwable getFailureCause()
public SimpleSlot getCurrentAssignedResource()
public InstanceConnectionInfo getCurrentAssignedResourceLocation()
public void setOperatorState(StateHandle operatorState)
public StateHandle getOperatorState()
public ExecutionGraph getExecutionGraph()
public void connectSource(int inputNumber,
IntermediateResult source,
JobEdge edge,
int consumerNumber)
public void setScheduleLocalOnly(boolean scheduleLocalOnly)
public boolean isScheduleLocalOnly()
public Iterable<Instance> getPreferredLocations()
null to indicate no location preference.public void resetForNewExecution()
public boolean scheduleForExecution(Scheduler scheduler, boolean queued) throws NoResourceAvailableException
NoResourceAvailableExceptionpublic void deployToSlot(SimpleSlot slot) throws JobException
JobExceptionpublic void cancel()
public void fail(Throwable t)
public void prepareForArchiving()
throws IllegalStateException
IllegalStateExceptionpublic void cachePartitionInfo(PartialInputChannelDeploymentDescriptor partitionInfo)
public String getSimpleName()
getTaskName(), 'x' is the parallel
subtask index as returned by getParallelSubtaskIndex()+ 1, and 'y' is the total
number of tasks, as returned by getTotalNumberOfParallelSubtasks().Copyright © 2014–2015 The Apache Software Foundation. All rights reserved.