org.apache.flink.api.common.JobID jobID
byte[] key
org.apache.flink.api.common.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
org.apache.flink.api.common.JobID jobID
JobVertexID vertexID
ExecutionAttemptID executionId
String taskName
int indexInSubtaskGroup
int numberOfSubtasks
org.apache.flink.configuration.Configuration jobConfiguration
org.apache.flink.configuration.Configuration taskConfiguration
String invokableClassName
List<E> producedPartitions
List<E> inputGates
int targetSlotNumber
List<E> requiredJarFiles
SerializedValue<T> operatorState
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
SerializedValue<T> operatorState
SerializableObject progressLock
org.apache.flink.api.common.JobID jobID
String jobName
org.apache.flink.configuration.Configuration jobConfiguration
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()].scala.concurrent.duration.FiniteDuration timeout
int numberOfRetriesLeft
long delayBeforeRetrying
boolean allowQueuedScheduling
ScheduleMode scheduleMode
boolean snapshotCheckpointsEnabled
JobStatus state
Throwable failureCause
int numFinishedJobVertices
Scheduler scheduler
ClassLoader userClassLoader
CheckpointCoordinator checkpointCoordinator
org.apache.flink.api.common.ExecutionConfig executionConfig
SerializableObject stateMonitor
ExecutionGraph graph
JobVertex jobVertex
ExecutionVertex[] taskVertices
IntermediateResult[] producedDataSets
List<E> inputs
int parallelism
boolean[] finishedSubtasks
int numSubtasksInFinalState
SlotSharingGroup slotSharingGroup
CoLocationGroup coLocationGroup
org.apache.flink.core.io.InputSplit[] inputSplits
List<E>[] inputSplitsPerSubtask
org.apache.flink.core.io.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
int numberOfCPUCores
long sizeOfPhysicalMemory
long sizeOfJvmHeap
long sizeOfManagedMemory
InetAddress inetAddress
int dataPort
String fqdnHostName
String hostName
Instance instance
boolean closed
InetSocketAddress address
int connectionIndex
SocketAddress address
ResultPartitionID partitionId
Throwable cause
IntermediateResultPartitionID partitionId
ExecutionAttemptID producerId
int expectedSequenceNumber
int actualSequenceNumber
String formatDescription
IntermediateDataSetID id
JobVertex producer
List<E> consumers
ResultPartitionType resultType
JobVertex target
DistributionPattern distributionPattern
IntermediateDataSet source
IntermediateDataSetID sourceId
Map<K,V> taskVertices
org.apache.flink.configuration.Configuration jobConfiguration
List<E> userJars
List<E> userJarBlobKeys
org.apache.flink.api.common.JobID jobID
String jobName
int numExecutionRetries
boolean allowQueuedScheduling
ScheduleMode scheduleMode
JobSnapshottingSettings snapshotSettings
JobVertexID id
ArrayList<E> results
ArrayList<E> inputs
int parallelism
org.apache.flink.configuration.Configuration configuration
String invokableClassName
org.apache.flink.core.io.InputSplitSource<T extends org.apache.flink.core.io.InputSplit> inputSplitSource
String name
SlotSharingGroup slotSharingGroup
CoLocationGroup coLocationGroup
String formatDescription
akka.actor.ActorRef jobmanager
akka.actor.ActorRef archive
scala.concurrent.duration.FiniteDuration timeout
File[] logDirs
org.apache.flink.configuration.Configuration configuration
akka.actor.ActorRef jobmanager
scala.concurrent.duration.FiniteDuration timeout
scala.Option<A> job
private Object readResolve()
scala.collection.Iterable<A> jobs
private Object readResolve()
org.apache.flink.api.common.JobID jobID
ExecutionGraph graph
private Object readResolve()
org.apache.flink.api.common.JobID jobID
private Object readResolve()
private Object readResolve()
private Object readResolve()
org.apache.flink.api.common.JobID jobID
JobVertexID vertexID
String taskName
int totalNumberOfSubTasks
int subtaskIndex
ExecutionAttemptID executionID
ExecutionState newExecutionState
long timestamp
String optionalMessage
private Object readResolve()
private Object readResolve()
JobGraph jobGraph
private Object readResolve()
JobGraph jobGraph
private Object readResolve()
org.apache.flink.api.common.JobID jobID
private Object readResolve()
org.apache.flink.api.common.JobID jobID
Throwable cause
private Object readResolve()
org.apache.flink.api.common.JobID jobID
private Object readResolve()
boolean success
scala.Option<A> error
private Object readResolve()
org.apache.flink.api.common.JobID jobID
JobStatus status
private Object readResolve()
org.apache.flink.api.common.JobID jobID
ExecutionGraph executionGraph
private Object readResolve()
private Object readResolve()
org.apache.flink.api.common.JobID jobID
private Object readResolve()
SerializedJobExecutionResult result
private Object readResolve()
byte[] splitData
private Object readResolve()
scala.collection.Iterable<A> taskManagers
private Object readResolve()
private Object readResolve()
org.apache.flink.api.common.JobID jobID
private Object readResolve()
private Object readResolve()
org.apache.flink.api.common.JobID jobID
private Object readResolve()
org.apache.flink.api.common.JobID jobID
JobVertexID vertexID
ExecutionAttemptID executionAttempt
private Object readResolve()
private Object readResolve()
org.apache.flink.api.common.JobID jobId
ResultPartitionID partitionId
ExecutionAttemptID taskExecutionId
IntermediateDataSetID taskResultId
private Object readResolve()
private Object readResolve()
private Object readResolve()
private Object readResolve()
InstanceID instanceID
private Object readResolve()
private Object readResolve()
scala.collection.Iterable<A> runningJobs
private Object readResolve()
scala.collection.Iterable<A> runningJobs
private Object readResolve()
org.apache.flink.api.common.JobID jobId
ResultPartitionID partitionId
private Object readResolve()
JobGraph jobGraph
boolean registerForEvents
private Object readResolve()
private Object readResolve()
String reason
private Object readResolve()
akka.actor.ActorRef jobManager
InstanceID instanceID
int blobPort
private Object readResolve()
akka.actor.ActorRef jobManager
InstanceID instanceID
int blobPort
private Object readResolve()
String reason
private Object readResolve()
akka.actor.ActorRef taskManager
InstanceConnectionInfo connectionInfo
HardwareDescription resources
int numberOfSlots
private Object readResolve()
String jobManagerAkkaURL
scala.concurrent.duration.FiniteDuration timeout
scala.Option<A> deadline
int attempt
private Object readResolve()
private Object readResolve()
InstanceID instanceID
byte[] metricsReport
private Object readResolve()
private Object readResolve()
private Object readResolve()
private Object readResolve()
private Object readResolve()
InstanceID instanceID
String stackTrace
private Object readResolve()
ExecutionAttemptID attemptID
private Object readResolve()
ExecutionAttemptID executionID
private Object readResolve()
ExecutionAttemptID executionID
Throwable cause
private Object readResolve()
ExecutionAttemptID taskExecutionId
IntermediateDataSetID taskResultId
IntermediateResultPartitionID partitionId
ExecutionState state
private Object readResolve()
TaskDeploymentDescriptor tasks
private Object readResolve()
ExecutionAttemptID executionID
private Object readResolve()
ExecutionAttemptID executionID
boolean success
String description
private Object readResolve()
TaskExecutionState taskExecutionState
private Object readResolve()
ExecutionAttemptID executionID
scala.collection.Seq<A> partitionInfos
private Object readResolve()
ExecutionAttemptID executionID
IntermediateDataSetID resultId
InputChannelDeploymentDescriptor partitionInfo
private Object readResolve()
org.apache.flink.api.common.JobID job
ExecutionAttemptID taskExecutionId
long checkpointId
SerializedValue<T> state
long timestamp
long timestamp
org.apache.flink.configuration.Configuration config
org.apache.flink.configuration.Configuration backingConfig
String prefix
private void writeObject(ObjectOutputStream oos) throws Exception
ExceptionString pathString
Serializable state
int numNetworkBuffers
int networkBufferSize
IOManager.IOMode ioMode
scala.Option<A> nettyConfig
scala.Tuple2<T1,T2> partitionRequestInitialAndMaxBackoff
private Object readResolve()
org.apache.flink.api.common.JobID jobID
ExecutionAttemptID executionId
ExecutionState executionState
byte[] serializedError
String[] tmpDirPaths
long cleanupInterval
scala.concurrent.duration.FiniteDuration timeout
scala.Option<A> maxRegistrationDuration
int numberOfSlots
org.apache.flink.configuration.Configuration configuration
private Object readResolve()
byte[] serializedData
int numberOfTaskManagers
int numberOfSlots
Copyright © 2014–2015 The Apache Software Foundation. All rights reserved.