private void readObject(ObjectInputStream in) throws IOException, ClassNotFoundException
IOExceptionClassNotFoundExceptionprivate void writeObject(ObjectOutputStream out) throws IOException
IOExceptionbyte[] key
JobID jobID
ResultPartitionID consumedPartitionId
ResultPartitionLocation consumedPartitionLocation
IntermediateDataSetID consumedResultId
int consumedSubpartitionIndex
DistributionPattern and the subtask indices of the producing and consuming task.InputChannelDeploymentDescriptor[] inputChannels
IntermediateDataSetID resultId
IntermediateResultPartitionID partitionId
ResultPartitionType partitionType
int numberOfSubpartitions
org.apache.flink.runtime.deployment.ResultPartitionLocation.LocationType locationType
ConnectionID connectionId
JobID jobID
JobVertexID vertexID
ExecutionAttemptID executionId
String taskName
int indexInSubtaskGroup
int numberOfSubtasks
Configuration jobConfiguration
Configuration taskConfiguration
String invokableClassName
List<E> producedPartitions
List<E> inputGates
int targetSlotNumber
List<E> requiredJarFiles
StateHandle operatorStates
long timestamp
long sequenceNumber
ExecutionVertex vertex
ExecutionAttemptID attemptId
long[] stateTimestamps
int attemptNumber
scala.concurrent.duration.FiniteDuration timeout
ConcurrentLinkedQueue<E> partialInputChannelDeploymentDescriptors
ExecutionState state
SimpleSlot assignedResource
Throwable failureCause
InstanceConnectionInfo assignedResourceLocation
StateHandle operatorState
JobID jobID
String jobName
Configuration jobConfiguration
ClassLoader userClassLoader
ConcurrentHashMap<K,V> tasks
List<E> verticesInCreationOrder
ConcurrentHashMap<K,V> intermediateResults
ConcurrentHashMap<K,V> currentExecutions
List<E> requiredJarFiles
List<E> jobStatusListenerActors
List<E> executionListenerActors
long[] stateTimestamps
System.currentTimeMillis() when
the execution graph transitioned into a certain state. The index into this array is the
ordinal of the enum value, i.e. the timestamp when the graph went into state "RUNNING" is
at stateTimestamps[RUNNING.ordinal()].Object progressLock
scala.concurrent.duration.FiniteDuration timeout
int numberOfRetriesLeft
long delayBeforeRetrying
boolean allowQueuedScheduling
ScheduleMode scheduleMode
JobStatus state
Throwable failureCause
Scheduler scheduler
int nextVertexToFinish
akka.actor.ActorContext parentContext
akka.actor.ActorRef stateCheckpointerActor
boolean checkpointingEnabled
long checkpointingInterval
Object stateMonitor
ExecutionGraph graph
AbstractJobVertex jobVertex
ExecutionVertex[] taskVertices
IntermediateResult[] producedDataSets
List<E> inputs
int parallelism
boolean[] finishedSubtasks
int numSubtasksInFinalState
SlotSharingGroup slotSharingGroup
CoLocationGroup coLocationGroup
InputSplit[] inputSplits
List<E>[] inputSplitsPerSubtask
InputSplitAssigner splitAssigner
ExecutionJobVertex jobVertex
Map<K,V> resultPartitions
ExecutionEdge[][] inputEdges
int subTaskIndex
List<E> priorExecutions
scala.concurrent.duration.FiniteDuration timeout
CoLocationConstraint locationConstraint
Execution currentExecution
List<E> locationConstraintInstances
boolean scheduleLocalOnly
StateHandle operatorState
int numberOfCPUCores
long sizeOfPhysicalMemory
long sizeOfJvmHeap
long sizeOfManagedMemory
InetAddress inetAddress
int dataPort
String fqdnHostName
String hostName
Instance instance
boolean closed
InetSocketAddress address
int connectionIndex
IntermediateResultPartitionID partitionId
ExecutionAttemptID producerId
int expectedSequenceNumber
int actualSequenceNumber
JobVertexID id
ArrayList<E> results
ArrayList<E> inputs
int parallelism
Configuration configuration
String invokableClassName
InputSplitSource<T extends InputSplit> inputSplitSource
String name
SlotSharingGroup slotSharingGroup
CoLocationGroup coLocationGroup
String formatDescription
IntermediateDataSetID id
AbstractJobVertex producer
List<E> consumers
ResultPartitionType resultType
AbstractJobVertex target
DistributionPattern distributionPattern
IntermediateDataSet source
IntermediateDataSetID sourceId
Map<K,V> taskVertices
Configuration jobConfiguration
List<E> userJars
List<E> userJarBlobKeys
JobID jobID
String jobName
int numExecutionRetries
boolean allowQueuedScheduling
ScheduleMode scheduleMode
JobGraph.JobType jobType
boolean checkpointingEnabled
long checkpointingInterval
String formatDescription
private Object readResolve()
private Object readResolve()
private Object readResolve()
String configDir
JobManagerMode executionMode
private Object readResolve()
AbstractID id
List<E> vertices
akka.actor.ActorRef jobmanager
akka.actor.ActorRef archive
scala.concurrent.duration.FiniteDuration timeout
File[] logDirs
Configuration configuration
akka.actor.ActorRef jobmanager
scala.concurrent.duration.FiniteDuration timeout
Configuration config
Configuration backingConfig
String prefix
int profilingInterval
int ioWaitCPU
int idleCPU
int userCPU
int systemCPU
int hardIrqCPU
int softIrqCPU
long totalMemory
long freeMemory
long bufferedMemory
long cachedMemory
long cachedSwapMemory
long receivedBytes
long transmittedBytes
JobID jobID
long profilingTimestamp
String instanceName
int userTime
int systemTime
int blockedTime
int waitedTime
JobVertexID vertexId
int subtask
ExecutionAttemptID executionId
int profilingInterval
byte[] state
Object stateObject
OperatorState<T> checkpointedState
int numNetworkBuffers
int networkBufferSize
IOManager.IOMode ioMode
scala.Option<A> nettyConfig
private Object readResolve()
JobID jobID
ExecutionAttemptID executionId
ExecutionState executionState
byte[] serializedError
String configDir
private Object readResolve()
String[] tmpDirPaths
long cleanupInterval
scala.concurrent.duration.FiniteDuration timeout
scala.Option<A> maxRegistrationDuration
int numberOfSlots
Configuration configuration
private Object readResolve()
int numberOfTaskManagers
int numberOfSlots
Copyright © 2014–2015 The Apache Software Foundation. All rights reserved.