org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator.forceNonParallel()方法的使用及代码示例

x33g5p2x  于2022-01-30 转载在 其他  
字(4.4k)|赞(0)|评价(0)|浏览(159)

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

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();

相关文章

微信公众号

最新文章

更多