Modifier and Type | Method and Description |
---|---|
OperatorStateBackend |
EmbeddedRocksDBStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
OperatorStateBackend |
RocksDBStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry)
Deprecated.
|
Modifier and Type | Class and Description |
---|---|
class |
DefaultOperatorStateBackend
Default implementation of OperatorStateStore that provides the ability to make snapshots.
|
Modifier and Type | Method and Description |
---|---|
abstract OperatorStateBackend |
AbstractStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
OperatorStateBackend |
StateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry)
Creates a new
OperatorStateBackend that can be used for storing operator state. |
Modifier and Type | Method and Description |
---|---|
OperatorStateBackend |
FsStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
OperatorStateBackend |
HashMapStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
Modifier and Type | Method and Description |
---|---|
OperatorStateBackend |
MemoryStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
protected Iterable<Tuple2<K,V>> |
BroadcastStateInputFormat.getElements(OperatorStateBackend restoredBackend) |
protected Iterable<OT> |
ListStateInputFormat.getElements(OperatorStateBackend restoredBackend) |
protected Iterable<OT> |
UnionStateInputFormat.getElements(OperatorStateBackend restoredBackend) |
Modifier and Type | Method and Description |
---|---|
OperatorStateBackend |
AbstractChangelogStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
Modifier and Type | Method and Description |
---|---|
OperatorStateBackend |
AbstractStreamOperatorV2.getOperatorStateBackend() |
OperatorStateBackend |
AbstractStreamOperator.getOperatorStateBackend() |
OperatorStateBackend |
StreamOperatorStateHandler.getOperatorStateBackend() |
OperatorStateBackend |
StreamOperatorStateContext.operatorStateBackend()
Returns the operator state backend for the stream operator.
|
protected OperatorStateBackend |
StreamTaskStateInitializerImpl.operatorStateBackend(String operatorIdentifierText,
PrioritizedOperatorSubtaskState prioritizedOperatorSubtaskStates,
CloseableRegistry backendCloseableRegistry) |
Modifier and Type | Method and Description |
---|---|
OperatorStateBackend |
BatchExecutionStateBackend.createOperatorStateBackend(Environment env,
String operatorIdentifier,
Collection<OperatorStateHandle> stateHandles,
CloseableRegistry cancelStreamRegistry) |
Constructor and Description |
---|
BeamDataStreamPythonFunctionRunner(String taskName,
ProcessPythonEnvironmentManager environmentManager,
String headOperatorFunctionUrn,
List<FlinkFnApi.UserDefinedDataStreamFunction> userDefinedDataStreamFunctions,
FlinkMetricContainer flinkMetricContainer,
KeyedStateBackend<?> keyedStateBackend,
OperatorStateBackend operatorStateBackend,
TypeSerializer<?> keySerializer,
TypeSerializer<?> namespaceSerializer,
TimerRegistration timerRegistration,
MemoryManager memoryManager,
double managedMemoryFraction,
FlinkFnApi.CoderInfoDescriptor inputCoderDescriptor,
FlinkFnApi.CoderInfoDescriptor outputCoderDescriptor,
FlinkFnApi.CoderInfoDescriptor timerCoderDescriptor,
Map<String,FlinkFnApi.CoderInfoDescriptor> sideOutputCoderDescriptors) |
BeamPythonFunctionRunner(String taskName,
ProcessPythonEnvironmentManager environmentManager,
FlinkMetricContainer flinkMetricContainer,
KeyedStateBackend<?> keyedStateBackend,
OperatorStateBackend operatorStateBackend,
TypeSerializer<?> keySerializer,
TypeSerializer<?> namespaceSerializer,
TimerRegistration timerRegistration,
MemoryManager memoryManager,
double managedMemoryFraction,
FlinkFnApi.CoderInfoDescriptor inputCoderDescriptor,
FlinkFnApi.CoderInfoDescriptor outputCoderDescriptor,
Map<String,FlinkFnApi.CoderInfoDescriptor> sideOutputCoderDescriptors) |
Modifier and Type | Method and Description |
---|---|
static BeamStateRequestHandler |
BeamStateRequestHandler.of(KeyedStateBackend<?> keyedStateBackend,
OperatorStateBackend operatorStateBackend,
TypeSerializer<?> keySerializer,
TypeSerializer<?> namespaceSerializer,
ReadableConfig config)
Create a
BeamStateRequestHandler . |
Constructor and Description |
---|
BeamOperatorStateStore(OperatorStateBackend operatorStateBackend) |
Modifier and Type | Method and Description |
---|---|
static void |
StreamingFunctionUtils.snapshotFunctionState(StateSnapshotContext context,
OperatorStateBackend backend,
Function userFunction) |
Copyright © 2014–2023 The Apache Software Foundation. All rights reserved.