【发布时间】:2015-11-17 19:32:25
【问题描述】:
由于性能测量,我想逐步执行为 Flink 编写的 Scala 程序,即
execute first operator; materialize result;
execute second operator; materialize result;
...
等等。原代码:
var filename = new String("<filename>")
var text = env.readTextFile(filename)
var counts = text.flatMap { _.toLowerCase.split("\\W+") }.map { (_, 1) }.groupBy(0).sum(1)
counts.writeAsText("file://result.txt", WriteMode.OVERWRITE)
env.execute()
所以我希望var counts = text.flatMap { _.toLowerCase.split("\\W+") }.map { (_, 1) }.groupBy(0).sum(1) 的执行是逐步的。
在每个操作员之后调用env.execute() 是正确的方法吗?
或者是在每次操作后写信给/dev/null,即调用counts.writeAsText("file:///home/username/dev/null", WriteMode.OVERWRITE),然后调用env.execute(),这是一个更好的选择? Flink 是否真的有类似 NullSink 的东西来实现这个目的?
编辑:我在集群上使用 Flink Scala Shell,并将应用程序设置为 parallelism=1 以执行上述代码。
【问题讨论】:
标签: scala apache-flink