异构Flink 集群中负载均衡算法研究与实现
2021-01-30汪志峰赵宇海王国仁
汪志峰 ,赵宇海* ,王国仁
(1.东北大学计算机科学与工程学院,沈阳,110169;2.北京理工大学计算机学院,北京,100081)
近年来,由于数据量呈指数式增长,处理大数据[1]的方式也发生了很大变化,迄今已经历三代引擎的改进:第一代以Hadoop[2]为代表,利用MapReduce[3]进行大数据处理;第二代以Spark[4]为代表,是基于RDD(基于内存计算)的微批流处理框架;第三代是以Flink 为代表的面向流,保证Exactly Once 的实时流处理框架.Flink[5]的计算平台可以实现毫秒级延迟下每秒处理上亿次的消息或者事件;同时,Flink 还提供一种Exactlyonce[6]的一致性语义,保证数据的正确性,使Flink大数据引擎可以提供金融级的数据处理能力.高吞吐、低延迟的性能使Flink 成为目前流处理的首选.与此同时,各个公司处理大数据的方式也发生了很大变化,例如阿里巴巴、滴滴出行、美团都大规模使用Flink 集群,阿里巴巴还开源自己的Blink 给Flink 社区,给Flink 带来了更好的性能优化以及方便的SQL 环境.如图1 所示的流程中,上游是事务处理、日志、点击事件等,经过Flink 的流处理到达下游;下游是处理之后的数据存储到数据库里或直接被应用所利用.Flink 在其中可以提供低延迟、高吞吐、Exactly Once 的处理.

图1 Flink 应用场景Fig.1 Flink application scenario
任务调度是Flink 很重要的功能.Flink 通过JobManger 任务调度器管理Slot,把任务分配到合适的Slot 等待TaskManager 执行.Flink 的任务调度图如图2 所示,Flink 集群启动后首先会启动一个JobManger 和一个或多个TaskManager,Client提交任务给JobManager,JobManager 再调度任务到各个TaskManager 去执行.在这个过程中TaskManager 把心跳和任务处理信息汇报给Job-Manager,TaskManager 之间以流的形式进行数据传输.在集群运行……
