org.apache.flink.runtime.execution.Environment.getTaskKvStateRegistry()方法的使用及代码示例

x33g5p2x  于2022-01-19 转载在 其他  
字(3.0k)|赞(0)|评价(0)|浏览(105)

本文整理了Java中org.apache.flink.runtime.execution.Environment.getTaskKvStateRegistry()方法的一些代码示例,展示了Environment.getTaskKvStateRegistry()的具体用法。这些代码示例主要来源于Github/Stackoverflow/Maven等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Environment.getTaskKvStateRegistry()方法的具体详情如下:
包路径:org.apache.flink.runtime.execution.Environment
类名称:Environment
方法名:getTaskKvStateRegistry

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),

相关文章