《基于 Apache Flink 的流处理》知识导图

原始笔记链接

关于本书

项目信息项目信息
书名《基于 Apache Flink 的流处理》作者[美] 比安‧霍斯克 / [美] 瓦西里基‧卡拉夫里
出版社中国电力出版社阅读日期2021年6月
豆瓣评分9.0我的评分☆☆☆☆☆

新手指引

安装及启动

brew install apache-flink
...
brew info apache-flink
...
cd /usr/local/Cellar/apache-flink/1.9.0/libexec/bin
./start-cluster.sh
访问: http://localhost:8081/
如果访问不同可以调整 libexec/conf/flink-conf.yaml 配置项

第一个 Demo

  1. 创建 demo 项目
mvn archetype:generate -DarchetypeCatalog=internal -DarchetypeGroupId=org.apache.flink -DarchetypeArtifactId=flink-quickstart-java -DarchetypeVersion=1.13.1
  1. Demo 程序

见下面 Demo 分析

  1. 监控 9000 端口
nc -l 9000
  1. 启动任务
./flink run -c bingjian.lbj.SocketTextStreamWordCount  /Users/bingjian.lbj/work/code/flink-demo/target/flink-demo-1.0-SNAPSHOT.jar 127.0.0.1 9000

此时可以在 web 页面上看到运行中的任务

  1. 在 nc 页面中输出
➜  flink-demo nc -l 9000
hello hello hello

可以在 log 中看到

/usr/local/Cellar/apache-flink/1.9.0/libexec/log

➜  log cat flink-bingjian.lbj-taskexecutor-1-C02YD3EGJHD2.local.out
(hello,1)
(hello,2)
(hello,3)

分析第一个 Demo

package bingjian.lbj;

import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.util.Collector;

/**
 * @author bingjian.lbj
 * @date 2021/6/14-10:57 上午
 **/
public class SocketTextStreamWordCount {

    public static void main(String[] args) throws Exception {
        //参数检查
        if (args.length != 2) {
            System.err.println("USAGE:\nSocketTextStreamWordCount <hostname> <port>");
            return;
        }

        String hostName = args[0];
        Integer port = Integer.parseInt(args[1]);

        //设置环境
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        //获取数据
        DataStream<String> text = env.socketTextStream(hostName, port);

        //计数
        DataStream<Tuple2<String, Integer>> counts = text.flatMap(new LineSplitter())
            .keyBy(0)
            .sum(1);
        counts.print();

优势

  • 同时支持事件时间(针对无序事件提供一致、精确的结果)和处理时间语义(能用在具有极低延迟需求的应用中)。

  • 提供精确一次(exactly-once)的状态一致性保障。

  • 毫秒级延迟

  • 层次化的 API 在表达能力和易用性方面各有权衡。

  • 用于最常见存储系统的连接器。

  • 支持高可用性配置(无单点失效)。

  • 允许在不丢失应用状态的前提下更新作业的程序代码,或进行跨 Flink 集群的作业迁移。

  • 提供详尽可定制的 metrics

  • Flink 也是一个成熟的批处理引擎。

数据交换策略

  • 转发策略

  • 广播策略

  • 基于键值的策略

  • 随机策略

窗口操作

  • 滚动窗口(tumbling window)将事件分配到长度固定且互不重叠的桶中。在窗口边界通过后,所有事件会发送给计算函数处理。基于数量的(count-based)滚动窗口定义了在触发计算前需要集齐多少条事件。

  • 滑动窗口(sliding window)将事件分配到大小固定且允许相互重叠的桶中,这意味着每个事件可能会同时属于多个桶。我们通过指定长度和滑动间隔来定义滑动窗口。

  • 会话窗口(session window)根据会话间隔(session gap)将事件分为不同的会话,该间隔值定义了会话在关闭前的非活动时间长度。

水位线:决定事件时间窗口的触发时机

全局进度指标,表示我们确信不会再有延迟事件到来的某个时间点,本质是一个逻辑时钟。 流处理系统很关键的一点是能提供某些机制来处理那些可能晚于水位线的迟到事件。根据应用需求的不同,你可能想直接忽略这些事件,将它们写入日志或利用它们去修正之前的结果。

支持有状态算子的挑战

  • 状态管理:系统需要有效地管理状态并保证它们不受并发更新的影响。

  • 状态划分:由于结果需要同时依赖状态和到来的事件,所以状态并行化会变得异常复杂。幸运的是,在很多情况下可以把状态按照键值划分,并独立管理每一部分。

  • 状态恢复:需要保证状态可以恢复,并且即使出现故障也要确保结果正确。

Flink 采用了轻量级检查点机制来实现精确一次结果保障

JobManager

  • 主进程,控制着单个应用程序的执行。每个应用都由一个不同的 JobManager 掌控。

  • JobManager 可以接受需要执行的应用,该应用会包含一个所谓的 JobGraph,即逻辑 Dataflow 图,以及一个大包了全部所需类、库以及其他资源的 Jar 文件。

  • JobManager 将 JobGraph 转化为名为 ExecutionGraph 的物理 Dataflow 图,该图包含了那些可以并行执行的任务。

  • JobManager 从 ResourceManager 申请执行任务的必要资源(TaskManager 处理槽)。一旦它收到了足够数量的 TaskManager 处理槽(slot),就会将 ExecutionGraph 中的任务分发给 TaskManager 来执行。

  • 在执行过程中,JobManager 还要负责所有需要集中协调的操作,如创建检查点。

ResourceManager

  • 负责管理 Flink 的处理资源单元——TaskManager 处理槽。当 JobManager 申请 TaskManager 处理槽时,ResourceManager 会指示一个拥有空闲处理槽的 TaskManager 将其处理槽提供给 JobManager。

  • 如果 ResourceManager 的处理槽数无法满足 JobManager 的请求,则 ResourceManager 可以和资源提供者通信,让它们提供额外容器来启动更多 TaskManager 进程。

  • ResourceManager 还负责终止空闲的 TaskManager 以释放计算资源。

TaskManager

  • Filnk 工作进程。

  • 每个 TaskManager 提供一定数量的处理槽。处理槽的数目限制了一个 TaskManager 可执行的任务数。TaskManager 在启动后,会向 ResourceManager 注册它的处理槽。当接收到 ResourceManager 的指示时,TaskManager 会向 JobManager 提供一个或多个处理槽。之后,JobManager 就可以向处理槽中分配任务来执行。

  • 执行期间,运行同一个应用不同任务的 TaskManager 之间会产生数据交换。

Dispatcher

  • 跨多个作业运行,提供了一个 REST 接口来让我们提交需要执行的应用.一旦某个应用提交执行,Dispatcher 会启动一个 JobManager 并将应用转交给它。

  • Dispatcher 同时还会启动一个 WebUI,用来提供有关作业执行的信息。

任务执行

TaskManager 允许同时执行多个任务

  • 数据并行(同一个算子)

  • 任务并行(不同算子)

  • 作业并行(不同应用)

TaskManager 会在同一个 JVM 进程内以多线程方式执行。

数据传输

基于信用值的流量控制

  • 接收任务会给发送任务授予一定的信用值,其实就是保留一些用来接收它数据的网络缓冲。一旦发送端收到信用通知,就会在信用值锁限定的范围内尽可能多地传输缓冲数据,并会附带上积压量(已经填满准备传输的网络缓冲数目)大小。

  • 接收端使用保留的缓冲来处理收到的数据,同时依据各发送端的积压量信息来计算所有相连的发送端在下一轮的信用优先级。

优势:应对数据倾斜。

Flink 检查点算法

  • 基于 Chandy-Lamport 分布式快照算法来实现的。不会暂定整个应用,而是会把生成检查点的过程和处理过程分离,这样在部分任务持久化状态的过程中,其他任务还可以继续执行。

  • 算法中会用到一类名为检查点分隔符(checkpoint barrier)的特殊记录。和水位线类似,这些检查点分隔符会通过数据源算子注入到常规的记录流中。相对其他记录,它们在流中的位置无法提前或延后。

  • 为了标识所属的检查点,每个检查点分隔符都会带有一个检查点编号,这样就把一条数据量从逻辑上分成了两个部分。

    • 所有先于分隔符的记录所引起的状态更改都会被包含在分隔符所对应的检查点之中;

    • 所有晚于分隔符的记录所引起的状态更改都会被纳入之后的检查点中。

Flink 保存点算法与 Connector 的 Checkpoint 基本一致

DataStream API

构建一个典型的 Flink 流式应用需要以下几步:

  1. 设置执行环境

    1. 根据上下文返回本地/远程环境:StreamExecutionEnvironment.getExecutionEnvironment()

    2. 创建一个本地的流式执行环境:StreamExecutionEnvironment.createLocalEnvironment()

    3. 创建一个远程的流式执行环境:StreamExecutionEnvironment.createRemoteEnvironment("host", 1234, "path/to/jarFile.jar")

    4. 可以使用 env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) 指定程序采用事件时间语义。

    5. 执行环境还提供了很多配置选项,例如设置程序并行度,启动容错等。

  2. 从数据源中读取一条或多条流

    1. StreamExecutionEnviroment 为我们提供了一系列创建流式数据源的方法。来源可以是消息队列,文件,可以是实时生成的。
  3. 通过一系列流式转换来实现应用逻辑

    1. 转换类型多种多样:有些会生成一个新的 DataStream(类型可能会发生变化);而另外一些不会修改 DataStream 中的记录,仅会通过分区或分组的方式将其重新组织。
  4. 选择性地将结果输出到一个或多个数据汇中

    1. 流式应用通常都会把结果发送到某些外部系统,例如 Kafka、文件系统或数据库。

    2. 有一些应用不会发出结果,而是将它们保存在内部,利用 Flink 的可查询式状态(queryable state)功能对外提供服务。

  5. 执行程序

    1. StreamExecutionEnviroment.execute()

转换操作

完成一个 DataStream API 程序的本质可以归结为:通过组合不同的转换来创建一个满足应用逻辑的 Dataflow 图.

DataStream API 转换分为四类:

基本转换:作用域单个事件

  • map() 可以指定 map 转换产生一个新的 DataStream。该转换将每个到来的事件传给一个用户自定义的映射器(user-defined mapper),后者针对每个输入只会返回一个输出事件。

  • filter() 利用一个作用在流中每条输入事件上的布尔条件来决定事件的去留。

  • flatMap() 类似于 map,但它可以对每个输入事件产生零个、一个或多个输出事件。

KeyedStream 转换:针对相同键值事件分组之后再进行处理

  • keyBy() 转换通过指定键值的方式将一个 DataStream 转化为一个 KeyedStream。流中的事件会根据各自键值被分到不同的分区,这样一来,有着相同键值的事件一定会在后续算子的同一个任务上处理。虽然键值不同的事件也可能会在同一个任务上处理,但任务函数所能访问的键值分区状态始终会被约束在当前键值的范围内。

  • 滚动聚合 作用于 KeyedStream 上,将生成一个包含聚合结果的 DataStream,例如 sum/min/max/minBy/maxBy

  • reduce() 是滚动聚合的泛化

多流转换: 将多条数据流合并为一条或将一条数据流拆分为多条流

  • union() 方法可以合并两条或多个类型相同的 DataStream,生成一个新的类型相同的 DataStream。这样后续的转换操作就可以对所有输入流中的元素统一处理。union 是 FIFO 的形式,无法保证顺序,不会对元素去重。

  • connect() 接收一个 DataStream 并返回一个 ConnectedStream 对象,该对象表示两个联结起来的流。

    • 提供 mapflatMap() 方法,分别接收一个 CoMapFunction 和 CoFlatMapFunction 作为参数。两个函数都是以两条输入流类型外加输出流类型作为其类型参数,它们为两条输入流定义了各自的处理方法。map1() 和 flatMap1() 用来处理第一条输入流事件,map2() 和 flatMap2() 用来处理第二条输入流事件。

    • 为了在 ConnectedStream s 上实现确定性的转换,connect() 可以与 keyBy() 和 broadcaset() 结合使用。

  • split()union 的逆操作,将输入流分割成两条或多条类型和输入流相同的输出流。split() 接收一个 OutputSelector 用来定义如何将数据流的元素分配到不同的命名输出中。

    • OutputSelector 中定义的 select() 方法会在每个输入事件到来时被调用,并随即返回一个 java.lang.Iterable[String] 对象。针对某记录所返回的一系列 String 值指定了该记录需要被发往哪些输出流。

    • split() 返回一个 SplitStream 对象,提供 select() 方法可以让我们通过指定输出名称的方式从 SplitStream 中选择一条或多条流。

分发转换:对流中的事件进行重新组织

  • Flink 各类分区转换对应了多种数据交换策略,这些操作定义了如何将事件分配给不同的任务。在构建程序时,系统会根据操作语义和配置的并行度自动选择数据分区策略并将数据转发到正确的目标。但是某些时候我们希望在应用级别控制这些分区策略,或者自定义分区器

  • shuffle() 随机数据交换策略

  • rebalance() 轮询均匀分配

  • rescale() 轮流方式分发,但分发目标仅限于部分猴急任务。当接收端任务远大于发送端任务时,改方法会更有效。

  • broadcast() 将所有事件复制并发送所有下游算子并行任务中

  • global 将所有事件发送下游算子第一个并行任务

  • 自定义 利用 partitionCustom() 方法自定义分区策略

设置并行度

  • 当提交一个 DataStream 到 JobManager 执行时,系统会生成一个 Dataflow 图并准备好用于执行的算子。每个算子都会产生或多个并发任务,每个任务负责处理算子的部分输入流。算子并行化任务的数目被称为该算子的并行度

  • 并行度可以在执行环境级别或单个算子级别控制

  • 默认为 CPU 数量

类型

  • Flink 利用类型信息的概念来表示数据类型,并且对于每种类型,都会为其生成特定的序列化器、反序列化器以及比较器。

支持的类型

  • 原始类型

  • Java(Tuple1…25)/Scala 元组

  • Scala 样例类

  • POJO(包括 APache Avro 生成的类)

  • 一些特殊类型

  • Flink 类型系统的核心类是 TypeInformation,它为系统生成序列化器和比较器提供了必要的信息。提供 TypeInformation 方式有两种

    • 实现 ResultTypeQueryable 接口扩展函数,在其中提供返回类型的 TypeInformation。

    • 定义 Dataflow 时使用 Java DataStream API 中的 returns() 方法来显示指定某算子的返回类型

定义键值和引用字段

  • 字段位置

  • 字段表达式,用于元组、POJO 和样例类

  • 键值选择器

实现函数

  • 函数类:以接口或抽象类的形式对外暴露,自己实现也同样如此。比如 MapFunction/FilterFunction/ProcessFunction

  • Lambda 函数

  • 富函数

    • 命名规则是 Rich 开头,后面跟着普通转换函数的名字,例如 RichMapFunction/RichFlatMapFunction

    • 使用富函数时,可以对应函数的生命周期实现两个额外的方法

      • open(): 初始化方法,在首次调用转换方法前会调用一次。

      • close(): 终止方法,在最后一次调用转换方法后调用一次。

    • 可以利用 getRuntimeContext() 方法访问函数的 RuntimeContext,并从中获取运行中信息。

基于时间和窗口的算子

配置时间特性: StreamExecutionEnvironment 可以接收以下值

  • ProcessTime: 指定算子根据处理机器的系统时钟决定数据流当前的时间。处理时间窗口基于机器时间出发,它可以覆盖触发时间点之前到达算子的任意元素。

  • EventTime: 指定算子根据数据自身包含的信息决定当前时间。每个事件时间都带有一个时间戳,而系统的逻辑时间是由水位时间线来定义。

    • 应用需要提供:每件事件都关联一个时间戳,该时间戳通常用来表示事件的实际发生时间。
  • IngestionTime: 指定每个接收的记录都把在数据源算子的处理时间作为事件时间的时间戳,并自动生成水位线。

处理函数

  • 除了基本功能它们还可以访问记录的时间戳和水位线,并支持注册在将来某个特定时间触发的计时器

  • 副输出功能还运行将记录发送到多个输出流中。处理函数常备用于构建事件驱动型应用,或实现一些内置窗口及转换无法实现的自定义逻辑。

8 类处理函数(ProcessFunction/KeyedProcessFunction/CoProcessFunction/ProcessJoinFunction/BroadcastProcessFunction/KeyedBroadcastProcessFunction/processWindowFunction/ProceessAllWindowFunction) 功能相似,适用于不同上下文,以 KeyedProcessFunction 举例。

  • 作用于 KeyedStream 上,针对流中的每条记录调用一次,并返回零个、一个或多个记录。

  • 所有处理函数都实现了 RichFunction 接口,因此支持 open()/close()/getRuntimeContext() 方法

  • KeyedProcessFunction[KEY, IN, OUT]提供两个方法:

    • processElement[v: IN, ctx: Context, out: Collector[out]]: 针对流中的每条记录都调用一次。可以像往常一样在方法中将结果记录传递给 Collector 发送出去。可通过 Context 访问时间戳、当前记录的键值以及 TimerService。此外,Context 还支持将结果发送到副输出。

    • onTimer(timestamp: Long, ctx: OnTimerContext, out: Collector[OUT]) 是一个回调函数,它会在之前注册的计时器触发时被调用。timestamp 参数给出了所触发计时器的时间戳,Collector 可用来发出记录。

时间服务和计时器

ContextOnTimerContext 对象中的 TimerService 提供了以下方法:

  • currentProcessingTime(): Long: 返回当前的处理时间

  • currentWatermark(): Long: 返回当前水位线的时间戳

  • registerProcessingTimeTimer(timestamp: Long): Unit: 针对当前键值注册一个处理时间计时器。当执行机器的处理时间达到给定的时间戳时,该计时器就会触发。

  • registerEventTimeTimer(timestamp: Long): Unit: 针对当前键值注册一个事件时间计时器。当更新后的水位线时间戳大于或等于计时器的时间戳时,就会触发。

  • deleteProcessingTimeTimer(timestamp: Long): Unit: 针对当前键值删除一个注册过的处理时间计时器。

  • deleteEventTimeTimer(timestamp: Long): Unit: 针对当前键值删除一个注册过的事件时间计时器。

向副输出发送数据

通过处理函数中的 Context 对象将记录发送到 OutputTag[X] 对象标识,其中 X 是副输出结果流的类型。

窗口算子

提供了一种基于有限大小的桶对事件进行分组,并对这些桶中的有限内容进行计算的方法。

定义窗口算子

用于键值分区可以并行计算,非键值分区的只能单线程处理

新建一个窗口算子需要指定两个窗口组件

  1. 一个用于决定输入流中的元素该如何划分的窗口分配器(window assigner)。会产生一个 WindowedStream(如果是用在非键值分区的 DataStream 上则是 AllWindowedStream)

  2. 一个作用于 WindowedStream(或 AllWindowedStream)上,用于处理分配到窗口中元素的窗口函数。

stream.keyBy(...)
    .window(...) // 指定窗口分配器
    .reduce/aggregate/process(...) // 指定窗口函数

内置窗口分配器

所有内置的窗口分配器都提供了一个默认的触发器,一旦(处理或事件)时间超过了窗口的结束时间就会触发窗口计算。

滚动窗口

会将元素放入大小固定且互不重叠的窗口中。

TumblingEventTimeWindows/TumblimgProcessingTimeWindows

val sensorData: DataStream[SensorReading] = ...
val avgTemp = sensorData.keyBy(_.id)
                // 将读数按照 1 秒事件时间窗口分组
                .window(TumblingEventTimeWindows.of(Time.seconds(1)))
                .process(new TemperatureAverager)

滑动窗口

将元素分配给大小固定且按指定滑动间隔移动的窗口。 如果滑动间隔小于窗口大小,则窗口会出现重叠,此时元素会被分配给多个窗口;如果滑动间隔大于窗口大小,则一些元素可能不会分配给任何窗口,因此可能会被直接丢弃。

会话窗口

将元素放入长度可变且不重叠的窗口中。会话窗口的边界由非活动间隔,即持续没有收到记录的时间间隔来定义。

在窗口上应用函数

可用于窗口的函数类型有两种:

  1. 增量聚合函数。它的应用场景是窗口内以状态形式存储某个值且需要根据每个加入窗口的元素对该值进行更新。此类函数通常会十分节省空间且最终会将聚合值作为单个结果发送出去。ReduceFunction 和 AggregateFunction 就属于增量聚合函数。

  2. 全量窗口函数。它会收集窗口内的所有元素,并在执行计算时对它们进行遍历。虽然全量窗口函数通常需要占用更多空间,但它和增量聚合函数相比,支持更复杂的逻辑。ProcessWindowFunction 就是一个全量窗口函数。

ReduceFunction

  • 会对分配给窗口的元素进行增量聚合,窗口只需要存储当前聚合结果,一个和 ReduceFunction 的输入及输出类型相同的值。

  • 每当收到一个新元素,算子都会以该元素和从窗口状态去除的当前聚合值为参数调用 ReduceFunction,随后会用 ReduceFunction 的结果替换窗口状态。

  • 每个窗口维护一个常量级别的小状态即可。

AggregateFunction

ProcessWindowFunction

自定义窗口算子

  • 可能需要更复杂的逻辑,例如:

    • 提前发出结果,对迟到的元素进行结果上的更新

    • 需要以特定记录作为开始或结束的边界

  • DataflowAPI 对外暴露了自定义窗口算子的接口和方法,可以实现自己的 分配器(assigner)触发器(trigger)以及移除器(evictor)。与上述的窗口函数协调工作。

窗口的生命周期

何时创建、包含了哪些信息以及何时删除。

窗口会在 WindowAssigner 首次向它分配元素时创建。因此,每个窗口至少会有一个元素。窗口内的状态由以下几部分组成:

  • 窗口内容:分配给窗口的元素,或当窗口算子配置了 ReduceFunction 或 AggregateFunction 时增量聚合所得的的结果。

  • 窗口对象:保存着用于区分窗口的信息。每个窗口对象都有一个结束时间戳,它定义了可以安全删除窗口及其状态的时间点。

  • 触发器计时器:可以在触发器中注册计时器,用于在将来某个时间点触发回调(例如对窗口进行计算或清理其内容)。这些计时器由窗口算子负责维护。

  • 触发器中的自定义状态:触发器可以定义和使用针对每个窗口、每个键值的自定义状态。

窗口算子会在窗口结束时间(由窗口对象中的结束时间戳定义)到达时删除窗口。当窗口需要删除时,窗口算子会自动清除窗口内容并丢弃窗口对。触发器状态和触发器中注册的计时器不会被清楚,因为这些状态对于窗口算子而言是不可见的。所以说,为了避免状态泄漏,触发器需要在 Trigger.clear() 方法中清除自身所有状态。

触发器

  • 用于定义何时对窗口进行计算并发出结果。它的触发条件可以是时间,也可以是某些特定的数据条件,如元素数量或某些观测到的元素值。

  • 触发器不仅能够访问时间属性和计时器,还可以使用状态,因此它在某种意义上等价于处理函数。

  • 自定义触发器还可以用来在水位线到达窗口的结束时间戳之前,为事件时间窗口计算并发出早期结果。这是一个在保守的水位线策略下依然可以产生(非完整的)低延迟结果的常用方法。

  • 每次调用触发器都会生成一个 TriggerResult,它用于决定窗口接下来的行为。有:

    • CONTINUE: 什么都不做

    • FIRE: 如果配置了 ProcessWindowFunction 则调用该函数发出结果;如果窗口只包含一个增量聚合函数,则直接发出当前聚合结果。窗口状态不会发生任何变化。

    • PURGE: 完全清除窗口内容,并删除窗口自身及其元数据。同时,调用 ProcessWindowFunction.clear() 方法来清理那些自定义的单个窗口状态。

    • FIRE_AND_PURGE: 先记下窗口计算(FIRE),随后删除所有状态及元数据(PURGE)。

基于时间的双流 Join

Dataflow API 中内置 有两个可以根据时间条件对数据流进行 Join 的算子:

  • 基于间隔的 Join

  • 基于窗口的 Join

基于间隔的 Join

  • 会对两条流中拥有相同键值以及彼此之间时间戳不超过某一指定间隔的时间进行 Join。

  • Join 成功的事件对会发送给 ProcessJoinFunction。下界和上界分别由负时间间隔和正时间间隔来定义,例如 between(Time.hour(-1), Time.minute(15))

  • 基于间隔的 Join 需要同时对双流的记录进行缓冲。

    • 对于第一个输入而言,所有时间戳大于当前数据线减去间隔上界的数据都会被缓冲起来;

    • 对于第二个输入而言,所有时间戳大于当前水位线加上间隔下界的数据都会被缓冲起来

  • 如果两条流的事件时间不同步,那么 Join 所需的存储就会显著增加,因为水位线总是由『较慢』的那条流来决定。

基于窗口的 Join

  • 将两条输入流中的元素分配到公共窗口中并在窗口完成时进行 Join(或 Cogroup)

  • 两条输入流都会根据各自的键值属性进行分区,公共窗口分配器会将二者的事件映射到公共窗口内(其实同时存储了两条流中数据)。当窗口的计时器触发时,算子会遍历两个输入中元素的每个组合(叉乘积)去调用 JoinFunction。

  • Join 与 Cogroup 总体逻辑相同,二者唯一区别是:Join 会为两侧输入中的每个事件对调用 JoinFunction;而 Cogroup 中用到的 GoGroupFunction 会以两个输入的元素遍历器为参数,只在每个窗口中被调用一次。

处理迟到数据

  • 丢弃迟到事件

  • 重定向迟到事件(到单独的数据流中)

  • 基于迟到事件更新结果

有状态算子和应用

  • 大部分复杂一些的操作都需要存储部分数据或中间结果。

  • 很多 Flink 内置的 DataStream 算子、数据源以及数据汇都是有状态的,它们需要对数据记录进行缓冲或者对中间结果及元数据加以维护

  • DataStream API 在用户自定义函数里暴露了状态的注册、维护及访问接口

  • 状态化流处理会在故障恢复、内存管理以及流式应用维护等很多方面对流处理引擎产生影响。

实现有状态函数

函数的状态类型有两种

  • 键值分区状态

  • 算子状态

Filnk 支持的三种算子状态:列表状态、联合列表状态以及广播状态。

在 RuntimeContext 中声明键值分区状态

  • 用户函数可以使用键值分区状态来存储和访问当前键值上下文中的状态。

  • 对于每个键值,Flink 都会维护一个状态实例。函数的键值分区状态实例会分布在函数所在算子的所有并行任务上。

  • 键值分区状态只能由作用于 KeyedStream 上面的函数使用。

  • Flinke 为键值分区状态提供了很多原语,定义了单个键值对应的状态结构:

    • ValueState[T]: 保存类型为 T 的单个值。

    • ListState[T]: 保存类型为 T 的元素列表。

    • MapState[K, V]: 保存一组键值映射。

    • ReducingState[T]: 提供和 ListState[T] 相同方法,但是其 add 会立即返回一个使用 ReduceFunction 聚合后的值。

    • AggregatingState[I, O]: 和 ReducingState 类似,但它使用了更加通用的 AggregateFunction 来聚合内部的值。

    • 所有原语都支持 State.clear() 方法进行清除。

通过 ListCheckpointed 接口实现算子列表状态

  • 若要在函数中使用算子列表状态,需要实现 ListCheckpointed 接口。该接口不像 ValueStateListState 那样直接在状态后端注册,而是需要将算子状态实现为成员变量并通过接口提供的回调函数与状态后端进行交互。

  • 提供两个方法:

    • List<T> snapshotState(Long checkpointId, Long timestamp): 以列表形式返回一个函数状态的快照。在 Flink 触发为有状态函数生成检查点时调用, checkpointId: 唯一且单调递增的检查点编号,timestamp: JobManager 开始创建检查点的机器时间戳。

    • void restoreState(List<T>: state): 根据提供的列表恢复函数状态。初始化函数状态时调用

使用联播的广播状态

在两条数据流上应用带有广播状态的函数需要三步:

  1. DataStream.broadcast() 创建一个 BroadcastStream 并提供一个或多个 MapStateDescriptor 对象。每个描述符都会为将来用于 BroadcastStream 的函数定义一个单独的广播状态。

  2. BroadcastStream 和一个 DataStreamKeyedStream 联结起来。必须将 BroadcastStream 作为参数传给 connect() 方法

  3. 在联结后的数据流上应用一个函数。根据另一条流是否已经按键值分区,该函数可能是 KeyedBroadcastProcessFunctionBroadcastProcessFunction

为有状态的应用开启故障恢复

  • Flink 为有状态的应用创建一致性检查点的机制:在所有算子都处理到应用输入流的某一个特定位置时,为全部内置或用户定义的有状态函数基于该时间点创建一个状态快照。

确定有状态应用的可维护性

关键两个参数:

  • 算子唯一标识

  • 最大并行度

有状态应用的性能及鲁棒性

选择状态后端

  • 状态后端负责存储每个状态实例的本地状态,并在生成检查点时将它们写入远程持久化存储。

  • 由于本地状态的维护及写入检查点的方式多种多样,所以状态后端被设计为『可插拔的』(pluggable),两个应用可以选择不同的状态后端实现来维护其状态。

  • 三种状态后端:

    • MemoryStateBackend

      • 存储在 TaskManager 进程的 JVM 堆中

      • 生成检查点时,MemoryStateBackend 会将状态发送至 JobManager 并保存到它的堆内存中。

      • 仅用于开发调试

    • FsStateBackend

      • 将本地状态保存在 TaskManager 的 JVM 中。

      • 写入远程持久化文件系统。

    • RocksDBStateBackend

      • 将全部状态存在本地 RocksDB 中

更新有状态应用

从保存点兼容性的角度来看,应用可以通过一下三种方式进行更新:

  1. 在不对已有状态进行更改或删除的前提下更新或扩展应用逻辑,包括向应用中添加有状态或无状态算子

  2. 从应用中移除某个状态

  3. 通过改变状态原语或数据类型来修改已有算子的状态

读写外部系统

应用的一致性保障

  • 应用若想提供端到端的精确一次性保障,需要一些特殊的数据汇连接器。根据情况不同,这些连接器可以使用两种技术来实现精确一次保障:

  • 幂等性写(idempotent write): 多次执行,但只会引起一次改变。

  • 事务性写(transactional write): 只有在上次成功的检查点之前计算的结果才会写入百步数据汇系统。

    • WAL(write ahead log,写前日志)

    • 2PC(two-phase commit, 二阶段提交)

实现自定义数据源函数

  • DataStream API 提供两个接口以及二者对应的 RichFunction 抽象类来实现数据源连接器

    • SourceFunction 和 RickSourceFunction 可用于定义非并行的数据源连接器,即只能以单任务运行。

    • ParallelSourceFunction 和 RichParallelSourceFunction 可用于定义能够同时运行多个任务实例的数据源连接器。

  • 两个关键方法

    • void run(SourceContext<T> ctx): 负责执行具体的记录读入或接收工作,Flink 单独开辟一个线程只会调用一次,通常会不断循环读取或接收数据将它们发出(无限流)

    • void cancel()

后续学习补充

  • Chandy-Lamport 分布式快照算法