本文整理了Java中org.apache.flink.runtime.execution.Environment.getInputGate()
方法的一些代码示例,展示了Environment.getInputGate()
的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。Environment.getInputGate()
方法的具体详情如下:
包路径:org.apache.flink.runtime.execution.Environment
类名称: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());
内容来源于网络,如有侵权,请联系作者删除!