When to Use a Distributed Dataflow Engine: Evaluating the Performance of Apache Flink

Ilya Verbitskiy, Lauritz Thamsen, Odej Kao · 2016

With the increasing amount of available data, distributed data processing systems like Apache Flink, Apache Spark have emerged that allow to analyze large-scale datasets. However, such engines introduce significant computational overhead compared to non-distributed implementations. Therefore, the question arises when using a distributed processing approach is actually beneficial. This paper helps to answer this question with an evaluation of the performance of the distributed data processing framework Apache Flink. In particular, we compare Apache Flink executed on up to 50 cluster nodes to single-threaded implementations executed on a typical laptop for three different benchmarks: TPC-H Query 10, Connected Components,, Gradient Descent. The evaluation shows that the performance of Apache Flink is highly problem dependent, varies from early outperformance in case of TPC-H Query 10 to slower runtimes in case of Connected Components. The reported results give hints for which problems, input sizes,, cluster resources using a distributed data processing system like Apache Flink or Apache Spark is sensible.

Read the paper · More papers on PaperTik