本文整理了Java中org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator.forceNonParallel()
方法的一些代码示例,展示了SingleOutputStreamOperator.forceNonParallel()
的具体用法。这些代码示例主要来源于Github
/Stackoverflow
/Maven
等平台,是从一些精选项目中提取出来的代码,具有较强的参考意义,能在一定程度帮忙到你。SingleOutputStreamOperator.forceNonParallel()
方法的具体详情如下:
包路径:org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator
类名称:SingleOutputStreamOperator
方法名:forceNonParallel
[英]Sets the parallelism and maximum parallelism of this operator to one. And mark this operator cannot set a non-1 degree of parallelism.
[中]
代码示例来源:origin: apache/flink
outTypeInfo,
operator
).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: apache/flink
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: com.data-artisans.streamingledger/da-streamingledger-runtime-serial
@Override
public ResultStreams translate(String name, List<InputAndSpec<?, ?>> streamLedgerSpecs) {
List<OutputTag<?>> sideOutputTags = createSideOutputTags(streamLedgerSpecs);
// the input stream is a union of different streams.
KeyedStream<TaggedElement, Boolean> input = union(streamLedgerSpecs)
.keyBy(unused -> true);
// main pipeline
String serialTransactorName = "SerialTransactor(" + name + ")";
SingleOutputStreamOperator<Void> resultStream = input
.process(new SerialTransactor(specs(streamLedgerSpecs), sideOutputTags))
.name(serialTransactorName)
.uid(serialTransactorName + "___SERIAL_TX")
.forceNonParallel()
.returns(Void.class);
// gather the sideOutputs.
Map<String, DataStream<?>> output = new HashMap<>();
for (OutputTag<?> outputTag : sideOutputTags) {
DataStream<?> rs = resultStream.getSideOutput(outputTag);
output.put(outputTag.getId(), rs);
}
return new ResultStreams(output);
}
}
代码示例来源:origin: dataArtisans/da-streamingledger
@Override
public ResultStreams translate(String name, List<InputAndSpec<?, ?>> streamLedgerSpecs) {
List<OutputTag<?>> sideOutputTags = createSideOutputTags(streamLedgerSpecs);
// the input stream is a union of different streams.
KeyedStream<TaggedElement, Boolean> input = union(streamLedgerSpecs)
.keyBy(unused -> true);
// main pipeline
String serialTransactorName = "SerialTransactor(" + name + ")";
SingleOutputStreamOperator<Void> resultStream = input
.process(new SerialTransactor(specs(streamLedgerSpecs), sideOutputTags))
.name(serialTransactorName)
.uid(serialTransactorName + "___SERIAL_TX")
.forceNonParallel()
.returns(Void.class);
// gather the sideOutputs.
Map<String, DataStream<?>> output = new HashMap<>();
for (OutputTag<?> outputTag : sideOutputTags) {
DataStream<?> rs = resultStream.getSideOutput(outputTag);
output.put(outputTag.getId(), rs);
}
return new ResultStreams(output);
}
}
代码示例来源:origin: org.apache.flink/flink-cep
comparator,
pattern.getAfterMatchSkipStrategy()
)).forceNonParallel();
代码示例来源:origin: org.apache.flink/flink-cep_2.11
comparator,
pattern.getAfterMatchSkipStrategy()
)).forceNonParallel();
代码示例来源:origin: org.apache.flink/flink-cep_2.10
nfaFactory,
false
)).forceNonParallel();
代码示例来源:origin: org.apache.flink/flink-streaming-java_2.10
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: org.apache.flink/flink-streaming-java
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11
return input.transform(opName, resultType, operator).forceNonParallel();
代码示例来源:origin: org.apache.flink/flink-streaming-java_2.11
return input.transform(opName, resultType, operator).forceNonParallel();
内容来源于网络,如有侵权,请联系作者删除!