【问题标题】:Apache Flink: stepwise executionApache Flink:逐步执行
【发布时间】: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


    【解决方案1】:

    Flink 默认使用流水线数据传输来提高作业执行的性能。但是,您也可以通过调用强制批量数据传输

    ExecutionEnvironment env = ...
    env.getConfig().setExecutionMode(ExecutionMode.BATCH_FORCED);
    

    这将分离两个运算符的执行(除非它们是链式的)。您可以从日志文件中获取每个任务的执行时间或查看 Web 仪表板。请注意,这不适用于链式运算符,即具有相同并行性且不需要网络洗牌的运算符。此外,您应该知道使用批量传输会增加程序的整体执行时间。我认为不可能真正分离流水线数据处理器中运算符的执行时间。

    在每个操作符不起作用后调用execute(),因为 Flink 还不支持将结果缓存到内存中。所以如果你执行operator 2,你要么需要将operator 1的结果写入某个持久化存储并再次读取它,要么再次执行operator 1。

    【讨论】:

    • 我将 flink 与单线程执行进行比较,因此我总是以 parallelism=1 来调用我的应用程序。那么是否可以分开执行?
    • 单线程执行对算子链有影响吗?
    • Chaining 是一种始终应用于 DataSet(批处理)程序的优化。只有在以下情况下才会发生链接:1) 两个运算符具有相同的并行度,2) 第二个运算符不需要分区,3) 第一个运算符只有一个后继。例如,具有相同并行度的两个映射运算符将始终链接。链接基本上意味着,数据不会序列化以进行传输,但生成的对象会立即转发到下一个函数。
    • 那么它在我的应用程序中怎么样?为了计算字数,我使用一张地图,然后使用一张减少(由 groupBy 和 sum 组成)。但是我还有三个运营商,对吧?在映射之后,需要对 reduce 进行分区。所以它不会被锁在地图之后?但是对于所有运算符,我的并行度默认为 1,对吗?然后上链?不是很明白...
    • 对于链接,需要满足所有条件。因此,在您的情况下, Mapper 和 Reducer 不会被链接,因为 Reduce 需要洗牌。然而,Flink 注入了一个连接到映射器的组合器。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2011-03-29
    • 2017-04-25
    • 2018-07-29
    • 1970-01-01
    • 2018-12-07
    相关资源
    最近更新 更多