Package | Description |
---|---|
org.apache.flink.contrib.streaming.state | |
org.apache.flink.runtime.checkpoint | |
org.apache.flink.runtime.deployment | |
org.apache.flink.runtime.execution | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.jobgraph.tasks | |
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. |
org.apache.flink.runtime.state | |
org.apache.flink.runtime.state.filesystem | |
org.apache.flink.runtime.state.memory | |
org.apache.flink.runtime.taskmanager | |
org.apache.flink.runtime.zookeeper | |
org.apache.flink.runtime.zookeeper.filesystem | |
org.apache.flink.streaming.runtime.operators |
This package contains the operators that perform the stream transformations.
|
org.apache.flink.streaming.runtime.tasks |
This package contains classes that realize streaming tasks.
|
Modifier and Type | Method and Description |
---|---|
<S extends Serializable> |
RocksDBStateBackend.checkpointStateSerializable(S state,
long checkpointID,
long timestamp) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
KeyGroupState.getKeyGroupState() |
SerializedValue<StateHandle<?>> |
SubtaskState.getState() |
Map<Integer,SerializedValue<StateHandle<?>>> |
TaskState.getUnwrappedKvStates(Set<Integer> keyGroupPartition)
Retrieve the set of key-value state key groups specified by the given key group partition set.
|
Modifier and Type | Method and Description |
---|---|
PendingCheckpoint.TaskAcknowledgeResult |
PendingCheckpoint.acknowledgeTask(ExecutionAttemptID executionAttemptId,
SerializedValue<StateHandle<?>> state,
long stateSize,
Map<Integer,SerializedValue<StateHandle<?>>> kvState) |
PendingCheckpoint.TaskAcknowledgeResult |
PendingCheckpoint.acknowledgeTask(ExecutionAttemptID executionAttemptId,
SerializedValue<StateHandle<?>> state,
long stateSize,
Map<Integer,SerializedValue<StateHandle<?>>> kvState) |
Constructor and Description |
---|
KeyGroupState(SerializedValue<StateHandle<?>> keyGroupState,
long stateSize,
long duration) |
SubtaskState(SerializedValue<StateHandle<?>> state,
long stateSize,
long duration) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
TaskDeploymentDescriptor.getOperatorState() |
Constructor and Description |
---|
TaskDeploymentDescriptor(SerializedValue<JobInformation> serializedJobInformation,
SerializedValue<TaskInformation> serializedTaskInformation,
ExecutionAttemptID executionAttemptId,
int subtaskIndex,
int attemptNumber,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors) |
Modifier and Type | Method and Description |
---|---|
void |
Environment.acknowledgeCheckpoint(long checkpointId,
StateHandle<?> state)
Confirms that the invokable has successfully completed all steps it needed to
to for the checkpoint with the give checkpoint-ID.
|
Modifier and Type | Method and Description |
---|---|
void |
Execution.setInitialState(SerializedValue<StateHandle<?>> initialState,
Map<Integer,SerializedValue<StateHandle<?>>> initialKvState) |
void |
Execution.setInitialState(SerializedValue<StateHandle<?>> initialState,
Map<Integer,SerializedValue<StateHandle<?>>> initialKvState) |
Modifier and Type | Interface and Description |
---|---|
interface |
StatefulTask<T extends StateHandle<?>>
This interface must be implemented by any invokable that has recoverable state and participates
in checkpointing.
|
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) |
Modifier and Type | Interface and Description |
---|---|
interface |
StreamStateHandle
A state handle that produces an input stream when resolved.
|
Modifier and Type | Class and Description |
---|---|
class |
AsynchronousStateHandle<T>
StateHandle that can asynchronously materialize the state that it represents. |
class |
LocalStateHandle<T extends Serializable>
A StateHandle that includes the operator states directly.
|
Modifier and Type | Method and Description |
---|---|
static <T extends StateHandle<?>> |
StateUtils.setOperatorState(StatefulTask<?> op,
StateHandle<?> state)
Utility method to define a common generic bound to be used for setting a
generic state handle on a generic state carrier.
|
Modifier and Type | Method and Description |
---|---|
abstract <S extends Serializable> |
AbstractStateBackend.checkpointStateSerializable(S state,
long checkpointID,
long timestamp)
Writes the given state into the checkpoint, and returns a handle that can retrieve the state back.
|
StateHandle<DataInputView> |
AbstractStateBackend.CheckpointStateOutputView.closeAndGetHandle()
Closes the stream and gets a state handle that can create a DataInputView.
|
abstract StateHandle<T> |
AsynchronousStateHandle.materialize()
Materializes the state held by this
AsynchronousStateHandle . |
<T extends Serializable> |
StreamStateHandle.toSerializableHandle()
Converts this stream state handle into a state handle that de-serializes
the stream into an object using Java's serialization mechanism.
|
Modifier and Type | Method and Description |
---|---|
static <T extends StateHandle<?>> |
StateUtils.setOperatorState(StatefulTask<?> op,
StateHandle<?> state)
Utility method to define a common generic bound to be used for setting a
generic state handle on a generic state carrier.
|
Modifier and Type | Class and Description |
---|---|
class |
FileSerializableStateHandle<T extends Serializable>
A state handle that points to state stored in a file via Java Serialization.
|
class |
FileStreamStateHandle
A state handle that points to state in a file system, accessible as an input stream.
|
Modifier and Type | Method and Description |
---|---|
<S extends Serializable> |
FsStateBackend.checkpointStateSerializable(S state,
long checkpointID,
long timestamp) |
<T extends Serializable> |
FileStreamStateHandle.toSerializableHandle() |
Modifier and Type | Class and Description |
---|---|
class |
ByteStreamStateHandle
A state handle that contains stream state in a byte array.
|
class |
SerializedStateHandle<T extends Serializable>
A state handle that represents its state in serialized form as bytes.
|
Modifier and Type | Method and Description |
---|---|
<S extends Serializable> |
MemoryStateBackend.checkpointStateSerializable(S state,
long checkpointID,
long timestamp)
Serialized the given state into bytes using Java serialization and creates a state handle that
can re-create that state.
|
<T extends Serializable> |
ByteStreamStateHandle.toSerializableHandle() |
Modifier and Type | Method and Description |
---|---|
void |
RuntimeEnvironment.acknowledgeCheckpoint(long checkpointId,
StateHandle<?> state) |
Constructor and Description |
---|
Task(JobInformation jobInformation,
TaskInformation taskInformation,
ExecutionAttemptID executionAttemptID,
int subtaskIndex,
int attemptNumber,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
MemoryManager memManager,
IOManager ioManager,
NetworkEnvironment networkEnvironment,
BroadcastVariableManager bcVarManager,
ActorGateway taskManagerActor,
ActorGateway jobManagerActor,
scala.concurrent.duration.FiniteDuration actorAskTimeout,
LibraryCacheManager libraryCache,
FileCache fileCache,
TaskManagerRuntimeInfo taskManagerConfig,
TaskMetricGroup metricGroup)
IMPORTANT: This constructor may not start any work that would need to
be undone in the case of a failing task deployment.
|
Modifier and Type | Method and Description |
---|---|
StateHandle<T> |
ZooKeeperStateHandleStore.add(String pathInZooKeeper,
T state)
Creates a state handle and stores it in ZooKeeper with create mode
CreateMode.PERSISTENT . |
StateHandle<T> |
ZooKeeperStateHandleStore.add(String pathInZooKeeper,
T state,
org.apache.zookeeper.CreateMode createMode)
Creates a state handle and stores it in ZooKeeper.
|
StateHandle<T> |
ZooKeeperStateHandleStore.get(String pathInZooKeeper)
Gets a state handle from ZooKeeper.
|
StateHandle<T> |
StateStorageHelper.store(T state)
Stores the given state and returns a state handle to it.
|
Modifier and Type | Method and Description |
---|---|
List<Tuple2<StateHandle<T>,String>> |
ZooKeeperStateHandleStore.getAll()
Gets all available state handles from ZooKeeper.
|
List<Tuple2<StateHandle<T>,String>> |
ZooKeeperStateHandleStore.getAllSortedByName()
Gets all available state handles from ZooKeeper sorted by name (ascending).
|
Modifier and Type | Method and Description |
---|---|
StateHandle<T> |
FileSystemStateStorageHelper.store(T state) |
Modifier and Type | Class and Description |
---|---|
static class |
GenericWriteAheadSink.ExactlyOnceState
This state is used to keep a list of all StateHandles (essentially references to past OperatorStates) that were
used since the last completed checkpoint.
|
Modifier and Type | Field and Description |
---|---|
protected TreeMap<Long,Tuple2<Long,StateHandle<DataInputView>>> |
GenericWriteAheadSink.ExactlyOnceState.pendingHandles |
Modifier and Type | Method and Description |
---|---|
TreeMap<Long,Tuple2<Long,StateHandle<DataInputView>>> |
GenericWriteAheadSink.ExactlyOnceState.getState(ClassLoader userCodeClassLoader) |
Modifier and Type | Class and Description |
---|---|
class |
StreamTaskStateList
List of task states for a chain of streaming tasks.
|
Modifier and Type | Method and Description |
---|---|
StateHandle<Serializable> |
StreamTaskState.getFunctionState() |
StateHandle<?> |
StreamTaskState.getOperatorState() |
Modifier and Type | Method and Description |
---|---|
void |
StreamTaskState.setFunctionState(StateHandle<Serializable> functionState) |
void |
StreamTaskState.setOperatorState(StateHandle<?> operatorState) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.