本文整理了Java中org.apache.flink.runtime.execution.Environment.getTaskKvStateRegistry()
方法的一些代码示例,展示了Environment.getTaskKvStateRegistry()
的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Environment.getTaskKvStateRegistry()
方法的具体详情如下:
包路径:org.apache.flink.runtime.execution.Environment
类名称:Environment
方法名:getTaskKvStateRegistry
[英]Returns the registry for InternalKvState instances.
[中]返回InternalKvState实例的注册表。
代码示例来源:origin: apache/flink
taskInfo.getMaxNumberOfParallelSubtasks(),
keyGroupRange,
environment.getTaskKvStateRegistry(),
TtlTimeProvider.DEFAULT,
metricGroup),
代码示例来源:origin: com.alibaba.blink/flink-runtime
@Override
public AbstractInternalStateBackend createInternalStateBackend(
Environment env,
String operatorIdentifier,
int numberOfGroups,
KeyGroupRange keyGroupRange) {
return new HeapInternalStateBackend(
numberOfGroups,
keyGroupRange,
env.getUserClassLoader(),
env.getTaskStateManager().createLocalRecoveryConfig(),
env.getTaskKvStateRegistry(),
isUsingAsynchronousSnapshots(),
env.getExecutionConfig()
);
}
代码示例来源:origin: com.alibaba.blink/flink-runtime
@Override
public AbstractInternalStateBackend createInternalStateBackend(
Environment env,
String operatorIdentifier,
int numberOfGroups,
KeyGroupRange keyGroupRange) {
return new HeapInternalStateBackend(
numberOfGroups,
keyGroupRange,
env.getUserClassLoader(),
env.getTaskStateManager().createLocalRecoveryConfig(),
env.getTaskKvStateRegistry(),
isUsingAsynchronousSnapshots(),
env.getExecutionConfig()
);
}
代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10
public <K> AbstractKeyedStateBackend<K> createKeyedStateBackend(
TypeSerializer<K> keySerializer,
int numberOfKeyGroups,
KeyGroupRange keyGroupRange) throws Exception {
if (keyedStateBackend != null) {
throw new RuntimeException("The keyed state backend can only be created once.");
}
String operatorIdentifier = createOperatorIdentifier(
headOperator,
configuration.getVertexID());
keyedStateBackend = stateBackend.createKeyedStateBackend(
getEnvironment(),
getEnvironment().getJobID(),
operatorIdentifier,
keySerializer,
numberOfKeyGroups,
keyGroupRange,
getEnvironment().getTaskKvStateRegistry());
// let keyed state backend participate in the operator lifecycle, i.e. make it responsive to cancelation
cancelables.registerClosable(keyedStateBackend);
// restore if we have some old state
Collection<KeyedStateHandle> restoreKeyedStateHandles =
restoreStateHandles == null ? null : restoreStateHandles.getManagedKeyedState();
keyedStateBackend.restore(restoreKeyedStateHandles);
@SuppressWarnings("unchecked")
AbstractKeyedStateBackend<K> typedBackend = (AbstractKeyedStateBackend<K>) keyedStateBackend;
return typedBackend;
}
代码示例来源:origin: org.apache.flink/flink-streaming-java
taskInfo.getMaxNumberOfParallelSubtasks(),
keyGroupRange,
environment.getTaskKvStateRegistry(),
TtlTimeProvider.DEFAULT,
metricGroup),
代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11
taskInfo.getMaxNumberOfParallelSubtasks(),
keyGroupRange,
environment.getTaskKvStateRegistry(),
TtlTimeProvider.DEFAULT,
metricGroup),
内容来源于网络,如有侵权,请联系作者删除!