Package | Description |
---|---|
org.apache.flink.api.common.accumulators | |
org.apache.flink.runtime.checkpoint | |
org.apache.flink.runtime.client | |
org.apache.flink.runtime.deployment | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.messages.accumulators | |
org.apache.flink.runtime.messages.checkpoint |
This package contains the messages that are sent between
JobManager
and TaskManager to coordinate the checkpoint snapshots of the
distributed dataflow. |
Modifier and Type | Method and Description |
---|---|
static Map<String,Object> |
AccumulatorHelper.deserializeAccumulators(Map<String,SerializedValue<Object>> serializedAccumulators,
ClassLoader loader)
Takes the serialized accumulator results and tries to deserialize them using the provided
class loader.
|
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
StateForTask.getState() |
Modifier and Type | Method and Description |
---|---|
boolean |
PendingCheckpoint.acknowledgeTask(ExecutionAttemptID attemptID,
SerializedValue<StateHandle<?>> state,
long stateSize) |
Constructor and Description |
---|
StateForTask(SerializedValue<StateHandle<?>> state,
long stateSize,
JobVertexID operatorId,
int subtask,
long duration) |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<Object>> |
SerializedJobExecutionResult.getSerializedAccumulatorResults() |
Constructor and Description |
---|
SerializedJobExecutionResult(JobID jobID,
long netRuntime,
Map<String,SerializedValue<Object>> accumulators)
Creates a new SerializedJobExecutionResult.
|
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
TaskDeploymentDescriptor.getOperatorState() |
Constructor and Description |
---|
TaskDeploymentDescriptor(JobID jobID,
JobVertexID vertexID,
ExecutionAttemptID executionId,
String taskName,
int indexInSubtaskGroup,
int numberOfSubtasks,
int attemptNumber,
Configuration jobConfiguration,
Configuration taskConfiguration,
String invokableClassName,
List<ResultPartitionDeploymentDescriptor> producedPartitions,
List<InputGateDeploymentDescriptor> inputGates,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
long recoveryTimestamp)
Constructs a task deployment descriptor.
|
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<Object>> |
ExecutionGraph.getAccumulatorsSerialized()
Gets a serialized accumulator map.
|
Modifier and Type | Method and Description |
---|---|
void |
Execution.setInitialState(SerializedValue<StateHandle<?>> initialState,
long recoveryTimestamp) |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<Object>> |
AccumulatorResultsFound.result() |
Constructor and Description |
---|
AccumulatorResultsFound(JobID jobID,
Map<String,SerializedValue<Object>> result) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
AcknowledgeCheckpoint.getState() |
Constructor and Description |
---|
AcknowledgeCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId,
SerializedValue<StateHandle<?>> state,
long stateSize) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.