大数据处理引擎 Apache Flink 梳理

图灵汇官网

Apache Flink 是一个开源计算平台,专门用于分布式数据流处理和批量数据处理。它能够在一个运行时环境中同时支持流处理和批处理两种类型的应用,这使其与其他开源方案有所不同。

现有的一些开源计算方案将流处理和批处理视为两种独立的应用类型,因为它们提供的服务级别协议(SLA)大相径庭:流处理通常需要低延迟和精确一次(Exactly-once)的处理保证,而批处理则需要高吞吐量和高效的处理能力。相比之下,Flink 从另一个角度看待这两种处理方式,将流处理视为输入数据流无界的应用,而将批处理视为输入数据流有界的应用。Flink 完全支持流处理,包括高吞吐、低延迟和高性能的处理能力,同时还支持多种窗口操作、有状态计算和容错机制。

Flink 的流处理特性包括高吞吐量、低延迟处理,支持事件时间窗口操作,保证精确一次的处理语义,支持灵活的窗口操作,支持反压功能,以及通过轻量级分布式快照实现的容错能力。Flink 还支持迭代计算和自动优化,避免不必要的昂贵操作,如 Shuffle 和排序。

Flink 的架构由三部分组成:客户端(Client)、任务管理器(TaskManager)和作业管理器(JobManager)。客户端用于提交用户任务,任务管理器执行具体的用户任务,而作业管理器则负责管理所有的任务管理器,并决定用户任务在哪些任务管理器上执行。Flink 提供的关键能力包括毫秒级别的低延迟处理、精确一次的处理保证、无单点故障的高可用性(HA)和水平扩展能力。

Flink 程序由流和转换这两个基本构建块组成。流是中间结果数据,而转换是对一个或多个输入流进行计算处理的操作。当一个 Flink 程序执行时,它会被映射为一个流数据流(Streaming Dataflow),这个流数据流由一系列流和转换操作符组成,类似于一个有向无环图(DAG)。一个流可以被分割成多个流分区,一个操作符可以被分割成多个操作符子任务,每个操作符子任务都在不同的线程中独立执行。

Flink 支持两种模式的数据流处理:一对一模式和重新分配模式。在一对一模式中,数据流的分区特性得以保持,而在重新分配模式中,数据流的分区会发生改变。紧密耦合的操作符可以被优化为操作符链,即多个操作符子任务被串连成一个操作符链,以提高执行效率。

处理流中的记录时,通常会包含三种典型的时间字段:事件时间、摄入时间和处理时间。Flink 使用水印(Watermark)来衡量时间,水印携带时间戳,并被插入到数据流中。水印有助于处理乱序的数据流,允许事件有一定的延迟。

窗口是流处理应用中的一种常见概念,用于收集最近一段时间内的数据并进行计算。窗口可以是时间驱动的,如每30秒一次,也可以是数据驱动的,如每100个元素一次。常见的窗口类型包括翻滚窗口、滑动窗口和会话窗口。

Flink 的容错机制基于分布式快照,通过周期性地插入屏障(barriers)到数据流中来实现。屏障将当前快照周期的数据与下一个周期的数据分隔开,确保在出现故障时能够恢复到故障前的状态。Flink 的快照机制借鉴了 Chandy-Lamport 算法,能够在多个任务管理器之间协调快照的生成和确认过程。

通过上述改写,我们希望保留并突出了原文中的关键信息点,同时尽量减少与原文的相似度。

本文来源: 图灵汇 文章作者: 雷科技