我正在尝试使用Apache Flink使用两种不同的算法处理数据流.我的伪代码如下:
env = getEnvironment();
DataStream<Event> inputStream = getInputStream();
// How to replicate the input stream?
Array[DataStream<Event>] inputStreams = inputStream.clone()
// apply different operations on the replicated streams
outputOne = inputStreams[0].map(func1);
outputTwo = inputStreams[1].map(func2);
...
outputOne.addSink(sink1);
outputTwo.addSink(sink2);
env.execute();
Run Code Online (Sandbox Code Playgroud)
我用Flink文档做了一些研究.似乎没有克隆流的概念.无论DataStream.iterate()也不DataStream.split()正在做的正是我想要的.是否有从源代码中多次创建流的替代方法?谢谢您的帮助.
"克隆"流非常简单,不需要专门的操作员.您可以在同一个上应用多个转换DataStream.所有下游转换都将使用完整的流.
所以在你的例子中你做了:
env = getEnvironment();
DataStream<Event> inputStream = getInputStream();
outputOne = inputStream.map(func1); // apply 1st transformation
outputTwo = inputStream.map(func2); // apply 2nd transformation
...
outputOne.addSink(sink1);
outputTwo.addSink(sink2);
env.execute();
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
964 次 |
| 最近记录: |