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

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

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

Environment.getJobID介绍

[英]Returns the ID of the job that the task belongs to.
[中]返回任务所属作业的ID。

代码示例

代码示例来源:origin: apache/flink

  1. this.jobId = env.getJobID();

代码示例来源:origin: apache/flink

  1. @Override
  2. public void init() throws Exception {
  3. final String iterationId = getConfiguration().getIterationId();
  4. if (iterationId == null || iterationId.length() == 0) {
  5. throw new Exception("Missing iteration ID in the task configuration");
  6. }
  7. final String brokerID = StreamIterationHead.createBrokerIdString(getEnvironment().getJobID(), iterationId,
  8. getEnvironment().getTaskInfo().getIndexOfThisSubtask());
  9. final long iterationWaitTime = getConfiguration().getIterationWaitTime();
  10. LOG.info("Iteration tail {} trying to acquire feedback queue under {}", getName(), brokerID);
  11. @SuppressWarnings("unchecked")
  12. BlockingQueue<StreamRecord<IN>> dataChannel =
  13. (BlockingQueue<StreamRecord<IN>>) BlockingQueueBroker.INSTANCE.get(brokerID);
  14. LOG.info("Iteration tail {} acquired feedback queue {}", getName(), brokerID);
  15. this.headOperator = new RecordPusher<>();
  16. this.headOperator.setup(this, getConfiguration(), new IterationTailOutput<>(dataChannel, iterationWaitTime));
  17. // call super.init() last because that needs this.headOperator to be set up
  18. super.init();
  19. }

代码示例来源:origin: apache/flink

  1. final String brokerID = createBrokerIdString(getEnvironment().getJobID(), iterationId ,
  2. getEnvironment().getTaskInfo().getIndexOfThisSubtask());

代码示例来源:origin: apache/flink

  1. () -> stateBackend.createKeyedStateBackend(
  2. environment,
  3. environment.getJobID(),
  4. operatorIdentifierText,
  5. keySerializer,

代码示例来源:origin: apache/flink

  1. checkpointStorage = stateBackend.createCheckpointStorage(getEnvironment().getJobID());

代码示例来源:origin: org.apache.flink/flink-statebackend-rocksdb_2.11

  1. this.jobId = env.getJobID();

代码示例来源:origin: org.apache.flink/flink-statebackend-rocksdb_2.10

  1. this.jobId = env.getJobID();

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

  1. public CheckpointStreamFactory createSavepointStreamFactory(StreamOperator<?> operator, String targetLocation) throws IOException {
  2. return stateBackend.createSavepointStreamFactory(
  3. getEnvironment().getJobID(),
  4. createOperatorIdentifier(operator, configuration.getVertexID()),
  5. targetLocation);
  6. }

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

  1. /**
  2. * This is only visible because
  3. * {@link org.apache.flink.streaming.runtime.operators.GenericWriteAheadSink} uses the
  4. * checkpoint stream factory to write write-ahead logs. <b>This should not be used for
  5. * anything else.</b>
  6. */
  7. public CheckpointStreamFactory createCheckpointStreamFactory(StreamOperator<?> operator) throws IOException {
  8. return stateBackend.createStreamFactory(
  9. getEnvironment().getJobID(),
  10. createOperatorIdentifier(operator, configuration.getVertexID()));
  11. }

代码示例来源:origin: org.apache.flink/flink-runtime_2.10

  1. public String brokerKey() {
  2. if (brokerKey == null) {
  3. int iterationId = config.getIterationId();
  4. brokerKey = getEnvironment().getJobID().toString() + '#' + iterationId + '#' +
  5. getEnvironment().getTaskInfo().getIndexOfThisSubtask();
  6. }
  7. return brokerKey;
  8. }

代码示例来源:origin: org.apache.flink/flink-runtime

  1. public String brokerKey() {
  2. if (brokerKey == null) {
  3. int iterationId = config.getIterationId();
  4. brokerKey = getEnvironment().getJobID().toString() + '#' + iterationId + '#' +
  5. getEnvironment().getTaskInfo().getIndexOfThisSubtask();
  6. }
  7. return brokerKey;
  8. }

代码示例来源:origin: org.apache.flink/flink-runtime_2.11

  1. public String brokerKey() {
  2. if (brokerKey == null) {
  3. int iterationId = config.getIterationId();
  4. brokerKey = getEnvironment().getJobID().toString() + '#' + iterationId + '#' +
  5. getEnvironment().getTaskInfo().getIndexOfThisSubtask();
  6. }
  7. return brokerKey;
  8. }

代码示例来源:origin: com.alibaba.blink/flink-runtime

  1. public String brokerKey() {
  2. if (brokerKey == null) {
  3. int iterationId = config.getIterationId();
  4. brokerKey = getEnvironment().getJobID().toString() + '#' + iterationId + '#' +
  5. getEnvironment().getTaskInfo().getIndexOfThisSubtask();
  6. }
  7. return brokerKey;
  8. }

代码示例来源:origin: org.apache.flink/flink-streaming-java

  1. @Override
  2. public void init() throws Exception {
  3. final String iterationId = getConfiguration().getIterationId();
  4. if (iterationId == null || iterationId.length() == 0) {
  5. throw new Exception("Missing iteration ID in the task configuration");
  6. }
  7. final String brokerID = StreamIterationHead.createBrokerIdString(getEnvironment().getJobID(), iterationId,
  8. getEnvironment().getTaskInfo().getIndexOfThisSubtask());
  9. final long iterationWaitTime = getConfiguration().getIterationWaitTime();
  10. LOG.info("Iteration tail {} trying to acquire feedback queue under {}", getName(), brokerID);
  11. @SuppressWarnings("unchecked")
  12. BlockingQueue<StreamRecord<IN>> dataChannel =
  13. (BlockingQueue<StreamRecord<IN>>) BlockingQueueBroker.INSTANCE.get(brokerID);
  14. LOG.info("Iteration tail {} acquired feedback queue {}", getName(), brokerID);
  15. this.headOperator = new RecordPusher<>();
  16. this.headOperator.setup(this, getConfiguration(), new IterationTailOutput<>(dataChannel, iterationWaitTime));
  17. // call super.init() last because that needs this.headOperator to be set up
  18. super.init();
  19. }

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11

  1. @Override
  2. public void init() throws Exception {
  3. final String iterationId = getConfiguration().getIterationId();
  4. if (iterationId == null || iterationId.length() == 0) {
  5. throw new Exception("Missing iteration ID in the task configuration");
  6. }
  7. final String brokerID = StreamIterationHead.createBrokerIdString(getEnvironment().getJobID(), iterationId,
  8. getEnvironment().getTaskInfo().getIndexOfThisSubtask());
  9. final long iterationWaitTime = getConfiguration().getIterationWaitTime();
  10. LOG.info("Iteration tail {} trying to acquire feedback queue under {}", getName(), brokerID);
  11. @SuppressWarnings("unchecked")
  12. BlockingQueue<StreamRecord<IN>> dataChannel =
  13. (BlockingQueue<StreamRecord<IN>>) BlockingQueueBroker.INSTANCE.get(brokerID);
  14. LOG.info("Iteration tail {} acquired feedback queue {}", getName(), brokerID);
  15. this.headOperator = new RecordPusher<>();
  16. this.headOperator.setup(this, getConfiguration(), new IterationTailOutput<>(dataChannel, iterationWaitTime));
  17. // call super.init() last because that needs this.headOperator to be set up
  18. super.init();
  19. }

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

  1. @Override
  2. public void init() throws Exception {
  3. final String iterationId = getConfiguration().getIterationId();
  4. if (iterationId == null || iterationId.length() == 0) {
  5. throw new Exception("Missing iteration ID in the task configuration");
  6. }
  7. final String brokerID = StreamIterationHead.createBrokerIdString(getEnvironment().getJobID(), iterationId,
  8. getEnvironment().getTaskInfo().getIndexOfThisSubtask());
  9. final long iterationWaitTime = getConfiguration().getIterationWaitTime();
  10. LOG.info("Iteration tail {} trying to acquire feedback queue under {}", getName(), brokerID);
  11. @SuppressWarnings("unchecked")
  12. BlockingQueue<StreamRecord<IN>> dataChannel =
  13. (BlockingQueue<StreamRecord<IN>>) BlockingQueueBroker.INSTANCE.get(brokerID);
  14. LOG.info("Iteration tail {} acquired feedback queue {}", getName(), brokerID);
  15. this.headOperator = new RecordPusher<>();
  16. this.headOperator.setup(this, getConfiguration(), new IterationTailOutput<>(dataChannel, iterationWaitTime));
  17. // call super.init() last because that needs this.headOperator to be set up
  18. super.init();
  19. }

代码示例来源:origin: org.apache.flink/flink-streaming-java

  1. final String brokerID = createBrokerIdString(getEnvironment().getJobID(), iterationId ,
  2. getEnvironment().getTaskInfo().getIndexOfThisSubtask());

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10

  1. public <K> AbstractKeyedStateBackend<K> createKeyedStateBackend(
  2. TypeSerializer<K> keySerializer,
  3. int numberOfKeyGroups,
  4. KeyGroupRange keyGroupRange) throws Exception {
  5. if (keyedStateBackend != null) {
  6. throw new RuntimeException("The keyed state backend can only be created once.");
  7. }
  8. String operatorIdentifier = createOperatorIdentifier(
  9. headOperator,
  10. configuration.getVertexID());
  11. keyedStateBackend = stateBackend.createKeyedStateBackend(
  12. getEnvironment(),
  13. getEnvironment().getJobID(),
  14. operatorIdentifier,
  15. keySerializer,
  16. numberOfKeyGroups,
  17. keyGroupRange,
  18. getEnvironment().getTaskKvStateRegistry());
  19. // let keyed state backend participate in the operator lifecycle, i.e. make it responsive to cancelation
  20. cancelables.registerClosable(keyedStateBackend);
  21. // restore if we have some old state
  22. Collection<KeyedStateHandle> restoreKeyedStateHandles =
  23. restoreStateHandles == null ? null : restoreStateHandles.getManagedKeyedState();
  24. keyedStateBackend.restore(restoreKeyedStateHandles);
  25. @SuppressWarnings("unchecked")
  26. AbstractKeyedStateBackend<K> typedBackend = (AbstractKeyedStateBackend<K>) keyedStateBackend;
  27. return typedBackend;
  28. }

代码示例来源:origin: org.apache.flink/flink-streaming-java

  1. () -> stateBackend.createKeyedStateBackend(
  2. environment,
  3. environment.getJobID(),
  4. operatorIdentifierText,
  5. keySerializer,

代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11

  1. () -> stateBackend.createKeyedStateBackend(
  2. environment,
  3. environment.getJobID(),
  4. operatorIdentifierText,
  5. keySerializer,

相关文章