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

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

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

Environment.getInputGate介绍

暂无

代码示例

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

  1. InputGate reader = getEnvironment().getInputGate(i);
  2. switch (inputType) {
  3. case 1:

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

  1. @Override
  2. public void invoke() throws Exception {
  3. RecordReader<SpeedTestRecord> reader = new RecordReader<>(
  4. getEnvironment().getInputGate(0),
  5. SpeedTestRecord.class,
  6. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  7. try {
  8. boolean isSlow = getTaskConfiguration().getBoolean(IS_SLOW_RECEIVER_CONFIG_KEY, false);
  9. int numRecords = 0;
  10. while (reader.next() != null) {
  11. if (isSlow && (numRecords++ % IS_SLOW_EVERY_NUM_RECORDS) == 0) {
  12. Thread.sleep(IS_SLOW_SLEEP_MS);
  13. }
  14. }
  15. }
  16. finally {
  17. reader.clearBuffers();
  18. }
  19. }
  20. }

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

  1. @Override
  2. public void invoke() throws Exception {
  3. RecordReader<SpeedTestRecord> reader = new RecordReader<>(
  4. getEnvironment().getInputGate(0),
  5. SpeedTestRecord.class,
  6. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  7. RecordWriter<SpeedTestRecord> writer = new RecordWriter<>(getEnvironment().getWriter(0));
  8. try {
  9. SpeedTestRecord record;
  10. while ((record = reader.next()) != null) {
  11. writer.emit(record);
  12. }
  13. }
  14. finally {
  15. reader.clearBuffers();
  16. writer.clearBuffers();
  17. writer.flushAll();
  18. }
  19. }
  20. }

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. InputGate reader = getEnvironment().getInputGate(i);
  2. switch (inputType) {
  3. case 1:

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories(),
  3. getEnvironment().getTaskManagerInfo().getConfiguration());
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. @Override
  2. public void invoke() throws Exception {
  3. this.headEventReader = new MutableRecordReader<>(
  4. getEnvironment().getInputGate(0),
  5. getEnvironment().getTaskManagerInfo().getTmpDirectories());

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

  1. @Override
  2. public void invoke() throws Exception {
  3. this.headEventReader = new MutableRecordReader<>(
  4. getEnvironment().getInputGate(0),
  5. getEnvironment().getTaskManagerInfo().getTmpDirectories());

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

  1. InputGate reader = getEnvironment().getInputGate(i);
  2. switch (inputType) {
  3. case 1:

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

  1. InputGate reader = getEnvironment().getInputGate(i);
  2. switch (inputType) {
  3. case 1:

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

  1. getEnvironment().getInputGate(currentReaderOffset),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories(),
  3. getEnvironment().getTaskManagerInfo().getConfiguration());
  4. readers[j] = getEnvironment().getInputGate(currentReaderOffset + j);

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

  1. getEnvironment().getInputGate(0),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){

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

  1. getEnvironment().getInputGate(0),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){

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

  1. getEnvironment().getInputGate(0),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories());
  3. } else if (groupSize > 1){

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

  1. getEnvironment().getInputGate(0),
  2. getEnvironment().getTaskManagerInfo().getTmpDirectories(),
  3. getEnvironment().getTaskManagerInfo().getConfiguration());

相关文章