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

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

本文整理了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

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

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

@Override
  public void invoke() throws Exception {
    RecordReader<SpeedTestRecord> reader = new RecordReader<>(
        getEnvironment().getInputGate(0),
        SpeedTestRecord.class,
        getEnvironment().getTaskManagerInfo().getTmpDirectories());
    try {
      boolean isSlow = getTaskConfiguration().getBoolean(IS_SLOW_RECEIVER_CONFIG_KEY, false);
      int numRecords = 0;
      while (reader.next() != null) {
        if (isSlow && (numRecords++ % IS_SLOW_EVERY_NUM_RECORDS) == 0) {
          Thread.sleep(IS_SLOW_SLEEP_MS);
        }
      }
    }
    finally {
      reader.clearBuffers();
    }
  }
}

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

@Override
  public void invoke() throws Exception {
    RecordReader<SpeedTestRecord> reader = new RecordReader<>(
        getEnvironment().getInputGate(0),
        SpeedTestRecord.class,
        getEnvironment().getTaskManagerInfo().getTmpDirectories());
    RecordWriter<SpeedTestRecord> writer = new RecordWriter<>(getEnvironment().getWriter(0));
    try {
      SpeedTestRecord record;
      while ((record = reader.next()) != null) {
        writer.emit(record);
      }
    }
    finally {
      reader.clearBuffers();
      writer.clearBuffers();
      writer.flushAll();
    }
  }
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

getEnvironment().getInputGate(0),
getEnvironment().getTaskManagerInfo().getTmpDirectories(),
getEnvironment().getTaskManagerInfo().getConfiguration());

相关文章