Adaptive Distributed Partitioning in Apache Flink
Theodoros Toliopoulos, Anastasios Gounaris · 2020
Dynamically adapting the workload of each worker in Flink is a challenging issue. In this work, we deal with a special case, where the data are conceptually split in contiguous overlapping regions. This scenario is encountered in several streaming applications, such as those employing nearest neighbor queries. We propose (i) an architecture for allowing such adaptations in Flink and (ii) specific data repartitioning techniques. We apply our proposal to a specific use case, namely continuous distancebased outlier detection. Our experimental evaluation provides insights into the efficiency and effectiveness of the approach.