【问题标题】:Is one Task one thread in Apache Flink是 Apache Flink 中的一个 Task 一个线程
【发布时间】:2020-05-15 02:59:49
【问题描述】:

我是 Flink 的新手。据我了解,在 Flink 中,一个 TaskManager 可以分为多个 slot,一个 slot 可以分配多个任务,一个任务是一个线程。

让我们看一下 WordCount 示例:

据我了解,一个任务就是一个线程,共有三个任务:Source + map()keyBy()/window()/apply()Sink。所以他们每个人都有自己的线程,这意味着我们在这个例子中需要三个线程。我们可以将三个任务(三个线程)放在一个槽中。

不过,现在我正在阅读它的官方文档:https://ci.apache.org/projects/flink/flink-docs-stable/dev/parallel.html

一个 Flink 程序由多个任务(转换/算子, 数据源和接收器)。 一个任务被分成几个并行的 执行实例,每个并行实例处理一个子集 任务的输入数据。一个任务的并行实例数 称为并行性。

如何理解“一个任务被分成多个并行执行实例”? “几个并行执行实例”是否意味着多线程?那么一个Task可以是多线程吗?

我现在很困惑。

【问题讨论】:

    标签: multithreading parallel-processing apache-flink flink-streaming


    【解决方案1】:

    措辞并不完美;任务有时在不同的上下文中具有不同的含义。

    在您的示例中,您展示了具有 3 个任务的程序的逻辑表示。由于它是一种逻辑表示,因此无法执行,因此考虑线程没有任何意义。

    当执行这样的逻辑表示时,它会被转换为物理表示。在最简单的情况下,每个逻辑任务都会产生 N 个物理任务,其中 N 是该任务的并行度。为了清楚起见,我们开始将物理任务称为子任务。

    可以粗略的说,每个子任务对应一个线程。但是,在算子链的情况下,子任务被合并到一个链中并在一个线程中执行。

    因此,在您的示例中,线程数由三个任务的并行度决定。所以你得到 N1+N2+N3 个线程。如果所有任务的并行度相同,则为 3*N。

    【讨论】:

    • 谢谢。现在很清楚了。还有一个问题:在这个例子中,如果我为子任务Source + map() 设置了两个并行度,这是否意味着两个线程将一起为同一个逻辑表示子任务Source + map() 工作?如果是这样,是否意味着我可以获得更好的性能,并且我不需要担心竞争条件等多线程问题?
    • 比如我设置两个度,说三个字来了,aaa,bbb,aaa。因此有可能一个线程持有aaabbb,而另一个线程持有最后一个aaa。我的意思是,第一个线程将标记aaa: 1 bbb: 1,第二个线程将标记aaa: 1,然后,Flink 将合并两个结果并标记aaa: 2 bbb: 1。我的理解对吗?
    • 如果任务在算子链中合并,它们只由一个物理子任务和线程处理。所以Source + map() 驻留在同一个线程中,并且在没有同步开销的情况下传递值。您的文字示例是正确的;合并是通过 keyBy 进行的,它将相等的键(在本例中为单词)混洗到同一个 subtask=thread。
    猜你喜欢
    • 1970-01-01
    • 2012-11-14
    • 2020-10-09
    • 1970-01-01
    • 2019-05-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多