Slot资源配置
Flink中有TaskSlot概念,每个taskSlot都有固定的资源,假设一个TaskManager有三个TaskSlots,那么每个TaskSlot会将TaskMananger中的内存均分,即每个任务槽的内存是总内存的1/3,分配资源意味着subtask不会与其他作业的subtask竞争内存,taskslot的作用就是分离任务的托管内存,不会发生cpu隔离。
用户无论是基于Standalone或者Yarn提交Flink任务时,都可以通过配置$FLINK_HOME/conf/flink-conf.yaml文件中的“taskmanager.numberOfTaskSlots”参数来指定每个JobManager启动后拥有几个taskslot。一个TaskManger可以配置成单Slot模式,这样这个JobManager上运行的任务就独占了整个JVM进程;一个TaskManager配置更多的taskSlot意味着更多的subtask可以共享同一个JVM,同一个JVM中的task共享TCP连接和心跳信息,但也意味着越多的taskslot争夺CPU资源。
关于每个TaskManager配置Slot个数的建议如下:
- 基于Standalone集群运行Flink任务时,建议每个TaskManager配置的slot个数与该节点Cpu core保持一致,这样能保证每个subtask尽量使用一个core来处理数据。
- 基于Yarn集群提交Flink任务时,会动态的申请TaskManager,每个TaskManager中的slot个数也是由参数taskmanager.numberOfTaskSlots决定,该参数默认为1。如果调整该值可以小于等于Yarn中单个Container最大可以申请的cpu核数(由参数yarn.scheduler.maximum-allocation-vcores决定,默认值为4),具体情况还需要结合每个subtask使用的内存决定该值设置多少,如果设置过大,可能每个subtask对应的内存会过小。
指定合适并行度
Flink中可以为每个算子设置不同并行度以应对不同业务处理逻辑达到数据处理最优化。设置高并行度往往会Flink数据处理的效率,但并行度也不是越大越好,太多并行度也会加重数据在多个Solt/TaskManaer之间数据传输压力,包括序列化和反序列化带来的压力,所以Flink任务设置合适的并行度很有必要。
确定任务合适并行度可以从如下几点考虑:
- 数据源并行度:如果Flink读取Source数据支持多并行度读取,那么可以设置并行度与Source数据源一致。例如Flink读取Kafka中数据,最好设置开始并行度与读取的Kafka Topic的分区数一致,每个分区都可以被一个task独立读取。如果后续数据处理数据的速度跟不上,也可以设置一个task读取topic多个分区数据,但数据源的并行度应该等于或大于Flink读取souce的并行度,不能出现Flink读取Source并行度大于Source源并行度,否则会出现一些并行度使用不上,浪费资源。
- 算子逻辑的复杂度:算子逻辑越复杂,相应的并行度就需要更高才能提供足够的吞吐量。
- 数据Sink并行度:Flink写出数据时,如果写出端支持多并行写出的数量固定,可以设置写出并行度与接收数据端并行保持一致。例如,Flink写出数据到Kafka,可以设置并行度为Kafka topic的分区数。如果写出支持并行,但没有固定数量,需要结合写出端承受数据写出压力来确定合适并行度。
- 系统资源可用性:一个Flink Application的并行度通常认为是所有算子中最大的并行度,在确定算子并行度时,应该考虑系统可用的资源,如果系统的CPU和内存资源有限,那么高并行度可能会导致任务竞争和低性能。
- 数据倾斜:数据倾斜是指数据分布不均衡,这可能会导致某些算子的并行度成为瓶领,在这种情况下,可以通过调整算子的并行度来解决数据倾斜的问题。
- 多次调整:最终的并行度设置往往需要通过多次调节来确定,每次调节后可以通过观察作业执行性能来一步步迭代调节确定并行度。
Flink中并行度可以从以下四个层面指定:
Operator Level (算子层面)
算子层面设置并行度是给每个算子设置并行度,直接在算子后面调用.setparallelism()方法,写入并行度即可,只是针对当前算子有效,注意一些算子不能设置并行度,例如:keyBy 返回的对象是KeyedStream,这种分组操作无法设置并行度,socketTextStream是非并行source,只支持1个并行度,也不能设置并行度。
// 算子层面设置并行度 ds.flatMap(line=>{line.split(" ")}).setParallelism(2)Execution Environment Level(执行环境层面)
执行环境层面设置并行度直接调用env.setParallelism()写入并行度即可,全局代码有效。
// 执行环境层面设置并行度 val env = StreamExecutionEnvironment.getExecutionEnvironment env.setParallelism(3)Client Level(客户端层面)
以上无论是算子层面还是执行环境层面设置并行度都会导致硬编码问题,修改并行度时不灵活,我们也可以在客户端提交Flink任务时通过指定命令参数-p来动态设置并行度,并行度作用于全局代码。
# 提交任务时通过 -p 参数指定并行度 ./flink run -m node1:8081 -p 4 -c com.wubaibao.flinkjava.code.chapter4.SocketWordCount /root/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar如果是基于WebUI提交任务,我们也可以基于WebUI指定并行度:

System Level(系统层面)
我们也可以直接在提交Flink任务的节点配置$FLINK_HOME/conf/flink-conf.yaml文件配置并行度,这个设置对于在客户端提交的所有任务有效,默认值为1。
#配置flink-conf.yaml文件 parallelism.default: 5
以上四种不同方式指定Flink并行度的优先级为:Operator Level>Execution Environment Level>Client Level>System Level,本地编写代码时如果没有指定并行度,默认的并行度是当前机器的cpu core数。
设置SlotSharingGroup共享组
默认情况下,Flink 允许 subtask 共享 taskSlot,即便它们是不同的 subtask,只要是来自于同一Flink作业即可(Flink不允许属于不同作业的task共享同一个slot),结果就是一个 slot 可以持有整个作业管道。

Flink中一个taskslot中可以运行多个subtask有什么好处呢?假设一个taskslot中只能运行一个subtask,上图中一共有13个subtask,对应的就需要13个slot资源,我们在提交Flink应用程序时需要关注我们程序中到底有多少subtask,然后再衡量Flink集群中slot个数是否足够,在一定程序上需要的slot资源较多。另外一个方面是在Flink中运行的task对CPU资源的占用不同,有CUP密集型task操作和CPU非密集型task操作情况,例如在Flink集群中source和map的操作只是读取数据进行转换,对应task运行占用的cpu资源极短,但是Window这种窗口聚合操作涉及大量数据计算,往往占用CPU资源时间长,这就会导致在运行任务时source/map、sink操作时间非常快,Window操作时间非常长,source/map对应的subtask会等待window对应的subtask执行,同样sink的对应的subtask也会等待window对应的subtask执行,站在集群slot角度上来看就出现了一些taskslot非常“繁忙”,一些taskslot非常“轻松”,集群的资源综合利用不高。
taskslot共享就可以很好地解决以上问题,Flink任务所有的subtask均衡的分散到不同的taskslot上执行,一个taskslot贯穿执行整个流程的subtask,这样每个taskslot、每个TaskManager上的资源使用情况非常均衡。所以允许 slot 共享有两个主要优点:
- Flink 集群所需的 taskSlot 和作业中使用的最大并行度恰好一样,不需要关注Flink程序总共包含多少个 subtask。
- 容易获得更好的资源利用。如果没有 slot 共享,非密集 subtask(source/map())将阻塞和密集型 subtask(window())一样多的资源。通过 slot 共享,确保繁重的 subtask 在 TaskManager 之间公平分配。
在Flink中实现taskslot共享是通过SlotSharingGroup(Slot共享组,简称SSG)实现的,默认在Flink中有名称为“default”的默认SSG,所有算子操作都在当前这个SSG中,所以我们在执行Flink代码时会自动进行slot组共享。我们也可以在代码中手动指定某些算子操作的SSG组做到某些操作独占一个slot,指定方式如下:
// 手动指定slotSharingGroup
someStream.filter(...).slotSharingGroup("name");
不显式指定SSG时所有算子操作使用的是default slot group 。显式指定后对应的算子操作使用的指定的slot group,只有指定同一个共享组的算子操作才会开启slot共享,不同slot group 的算子操作是分配到不同的slot上执行的,如果一个Flink 任务有多个共享组,那么该Flink任务所需的总slot个数就是每个共享组最大并行度的总和。
使用细粒度资源管理
在Flink架构中,TaskManager中会划分多个Slot资源,Slot是Flink运行时进行资源调度和资源分配的基本单元。

之前的Flink版本中,资源请求只包含所需的Slot,TaskManager有固定数量且资源相同Slot来满足用户资源请求,相当于是粗粒度的资源管理,现在Flink支持细粒度的资源管理,通过细粒度的资源管理,用户可以指定资源配置来对Slot进行请求,Flink根据用户的资源配置从TaskManager中动态剪切一个完全匹配的Slot,如上图所示,需要一个具有0.25 Core和1GB内存的Slot,Flink为其分配Slot 1。
注意:对于用户没有指定资源配置的资源请求,Flink会自动决定资源配置,目前默认的资源配置是根据TaskManager总资源和TaskManager.numberOfTaskSlots计算的,相当于是粗粒度资源管理。如上图所示,TaskManager的总资源为1Core和4G内存,当前TaskManager的Slot数量设置为2,那么每个Slot将会有0.5个core和2G内存。
在上图右侧图中是细粒度资源配置,TaskManager分配Slot 1和Slot 2后,TaskManager中剩余的可用内存为0.25 Core和1G内存,这些空闲资源可以进一步划分,以满足其他资源需求。
在Flink1.14版本中提出的细粒度资源调度是基于SlotSharingGroup的资源配置接口来实现,可以为任务中的每个SSG指定不同的资源可以最大化资源资源利用效率。
要使用细粒度资源管理,需要做以下操作:
配置启用细粒度资源管理
在flink-conf.yaml配置文件中配置 cluster.fine-grained-resource-management.enabled 为 true,早期Flink版本中没有此配置,如果配置上会有异常报错。
代码中指定资源需求
在代码中通过创建Slot Sharing Groups(Slot共享组)定义了细粒度的资源需求。在Flink内部SlotSharingGroup会告诉JobManager 哪些operator/tasks可以放在同一个Slot中。关于在代码中定义SlotSharingGroup和哪些算子使用对应的SSG,有以下两种方式:
构建SlotSharingGroup对象实例并指定资源,通过slotSharingGroup(String name)方式附加到算子上。
这种方式在创建SlotSharingGroup时指定共享组所需的资源,然后给算子通过 slotSharingGroup(String name) 方式来设置Slot共享组的名称,但是最后需要通过StreamExecutionEnvironment.registerSlotSharingGroup(SlotSharingGroup ssg) 注册这些SSG对象。
//创建SSG共享组对象指定资源配置 SlotSharingGroup ssgA = SlotSharingGroup.newBuilder("a") .setCpuCores(1.0) .setTaskHeapMemoryMB(10) .build(); //指定构建SSG共享组对象的字符串名称“a” ,使用当前SSG共享组资源,后续需要注册 someStream.filter(...).slotSharingGroup("a") //通过env注册名称为“a”的SSG共享组对象 env.registerSlotSharingGroup(ssgA);构建SlotSharingGroup对象实例并指定资源,通过slotshareinggroup (SlotSharingGroup ssg)附加到算子上。
同样,这种方式也是在创建SlotSharingGroup对象时指定SSG共享组的资源情况,给算子指定SSG共享组时直接通过slotshareinggroup (SlotSharingGroup ssg)即可。
//创建SSG共享组对象指定资源配置 SlotSharingGroup ssgB = SlotSharingGroup.newBuilder("b") .setCpuCores(0.5) .setTaskHeapMemoryMB(10) .build(); //直接指定SSG共享组对象名称来使用SSG共享组资源 DataStream<...> ds1 = someStream.filter(...).slotSharingGroup(ssgB)注意:无论以上使用那种方式指定SSG共享组资源,每个SSG共享组只能附加到一个指定的资源,任何冲突都将导致作业编译失败。此外在构造SlotSharingGroup (Slot共享组)实例时,可以为Slot共享组设置以下资源信息:
资源 解释 CPU cores 必须项,定义需要多少个 CPU 内核,需要显式配置正值。 Task Heap Memory 必须项,定义需要多少 Task 堆内存,需要显式配置正值。 Task Off-Heap Memory 定义需要多少 Task 堆外内存,可以是 0。Managed Memory 定义需要多少任务托管内存,可以是 0。External Resources 定义所需的外部资源,可以是空的。
Flink使用异步IO
Flink的异步I/O是一个非常受欢迎的特性,由阿里巴巴贡献给社区,并在1.2版本中引入,它的主要目的是解决与外部系统交互时网络延迟成为系统瓶颈的问题,外部系统往往是外部数据库。
在Flink流计算系统中,与外部数据库进行交互是常见的需求,通常情况下,我们会发送一个查询请求到数据库并等待结果返回,这期间无法发送其他请求,这种同步访问方式会导致阻塞,阻碍了吞吐量和延迟,为了解决这个问题,引入了异步模式,能够并发地处理多个到外部数据库的请求,下图为官方提供的异步IO原理图。

在Flink中使用异步I/O,我们可以连续发送多个查询请求到数据库,并在回复返回时处理每个回复,而不需要阻塞等待,这种并发处理的方式极大地减少了延迟。
异步I/O专门用于解决Flink计算过程中与外部系统的交互问题,特别需要注意的是为了提高Flink与外部系统交互能力,也可以提高Flink的并行度进而提高Flink处理数据的吞吐量,这种方式会付出更高的资源成本,如:更多的task、更多的内存缓存、更高的网络连接,而异步IO方式相对这种方式是基于某一个task之上的扩展,重复利用一个task资源做更多的事情,提高了Flink性能和资源利用率。
在Flink中查询外界数据库数据时要使用异步IO需要满足如下条件之一:
- 数据库(或K/V存储系统)提供支持异步请求的客户端,例如Java的Vertx。
- 对于不支持异步请求客户端的外部系统可以使用线程池模拟异步客户端。
大状态中设置TTL
Flink程序运行时随着时间推移,Flink状态大小也会持续增长,默认状态存储在内存中,时间久了就会给内存带来很大压力,我们可以调用.clear()方法直接清除状态,如果业务逻辑不允许使用clear()方法直接清除状态,我们可以通过配置状态的”生存时间”(Time-to-Live,TTL)来限制状态在内存中存在的时间,当状态存在的时间超过设置的TTL时,系统将自动尽快清除该状态,避免状态一直占用内存空间。
在代码中给键控状态设置TTL和应用TTL方式如下:
import org.apache.flink.api.common.state.StateTtlConfig;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.time.Time;
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.seconds(1))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
ValueStateDescriptor<String> stateDescriptor = new ValueStateDescriptor<>("text state", String.class);
stateDescriptor.enableTimeToLive(ttlConfig);
以上设置TTL的方法解释如下:
newBuilder(…)方法
该方法必须指定,需要传入一个时间,通过该方法来指定TTL状态生存时间。生存时间TTL计时是以Flink系统处理时间为基础,不支持事件事件。
setUpdateType(…)方法
该方法可选,通过该方法来设置何时更新状态失效时间。默认值为StateTtlConfig.UpdateType.OnCreateAndWrite,表示仅在创建和写入状态时更新TTL。还可以设置为StateTtlConfig.UpdateType.OnReadAndWrite,表示所有读与写状态时更新TTL,这种方式只要对状态在TTL时间内进行读取,那么该状态生存时间就会一直更新延后。
setStateVisibility(…)方法
该方法可选,该方法设置状态的可见性。默认值为StateTtlConfig.StateVisibility.NeverReturnExpired,表示状态数据过期就不会返回。还可以设置为StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp,表示状态数据即使过期,只要Flink还没有清除该状态就返回。
下面以ValueState为例来测试状态生存时间TTL。
需求:读取Socket中基站通话数据,统计每个主叫通话总时长。
Java代码实现
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
/**
* Socket中数据如下:
* 001,186,187,busy,1000,10
* 002,187,186,fail,2000,20
*/
DataStreamSource<String> ds = env.socketTextStream("node5", 9999);
//对ds进行转换处理,得到StationLog对象
SingleOutputStreamOperator<StationLog> stationLogDS = ds.map(new MapFunction<String, StationLog>() {
@Override
public StationLog map(String line) throws Exception {
String[] arr = line.split(",");
return new StationLog(
arr[0].trim(),
arr[1].trim(),
arr[2].trim(),
arr[3].trim(),
Long.valueOf(arr[4]),
Long.valueOf(arr[5])
);
}
});
stationLogDS.keyBy(stationLog -> stationLog.callOut)
.map(new RichMapFunction<StationLog, String>() {
//定义ValueState 用来存放同一个主叫号码的通话总时长
private ValueState<Long> valueState;
@Override
public void open(Configuration parameters) throws Exception {
//定义状态TTL
StateTtlConfig ttlConfig = StateTtlConfig
//设置状态有效期为10秒
.newBuilder(org.apache.flink.api.common.time.Time.seconds(10))
//设置状态更新类型为OnCreateAndWrite,即状态TTL创建和写入时更新
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
//设置状态可见性为NeverReturnExpired,即状态过期后不返回
设置barrier对齐和非对齐
当Flink应用程序有多个并行度或者Flink上下游算子并行度不一致时,barrier上下游传递时涉及到barrier广播和barrier对齐机制。当上游数据向下游多个并行度中发送barrier时,需要对barrier进行广播,保证下游各个并行度barrier一致;当上游多个并行度向下游少量并行度传递barrier时,需要对brrier进行对齐,对齐是指下游每个并行度都要等到相同的barrier到达时才能进行 snapshot 快照状态的保存。
下图是Flink中barrier对齐机制的示意图:

在barrier对齐机制中,下游barrier先到达分区会等待barrier未到达的分区以达到barrier对齐目的,这个对齐过程中就会涉及barrier先到达分区中数据的缓存,如果多个并行度中处理数据的速度不一致会导致下游任务堆积大量缓存数据,可能会造成Flink内存和磁盘负载压力,同时也使Flink 整体Checkpoint的时间延后变长,此外,在Flink中当数据流处理不过来时还会有反压机制(反压机制是指控制数据源的产生速度,避免数据积压),反压机制会限制数据流的流动,导致barrier在一些并行度中流速变慢,这样更近一步导致Checkpoint的时间延后的更长,导致恶性循环。
为了解决以上barrier对齐机制可能带来的问题,在Flink1.11后引入了barrier不对齐机制。下图是barrier不对齐机制的示意图:

当流速快的barrier到达下游算子的input buffer后,Flink会将该barrier插入到该下游算子的output buffer的最前面,并将该barrier发送给后续的算子。同时当前算子会对自身进行checkpoint快照,包括当前的状态以及所有input buffers、output buffers以及流速慢的barrier之前的数据都会保存到状态后端中(注意:在进行checkpoint快照时,流速慢的barrier会被移除,并不会继续流动下去),这样当Flink应用程序异常中断恢复到此次checkpoint时,未计算之前的状态、barrier不对齐对应的input buffers、output buffers 数据会重新恢复到各个流中并保证数据的一致性和准确。值得注意的是barrier不对齐机制中需要向状态中保持更多的数据。
通过上文对Flink barrier对齐和不对齐机制的了解,我们发现两者各有优缺点:
barrier对齐机制
优点:状态后端需要保存的数据少。
缺点:缓存堆积数据、Flink内存和磁盘负载有压力、checkpoint时间延长。
barrier不对齐机制
优点:多并行度中只要有一个并行度中barrier到达,就会触发checkpoint,加快checkpoint进行,不容易出现数据反压问题。
缺点:状态后端保存数据多,状态恢复时比较慢。
在Flink中对于简单数据处理作业建议使用轻量级的barrier对齐机制,对于一些计算复杂导致任务出现数据高反压、checkpoint超时难以完成的的作业场景建议使用barrier不对齐机制,这样可以加快checkpoint进行、有效缓解数据高反压带来的一系列连锁问题。
代码中设置uid
为了能够在作业的不同版本之间以及Flink的不同版本之间顺利升级,强烈推荐程序员通过手动给算子赋予ID,这些ID将用于确定每一个算子的状态范围。如果不手动给各算子指定ID,则会由Flink自动给每个算子生成一个ID。而这些自动生成的ID依赖于程序的结构,并且对代码的更改是很敏感的。因此,强烈建议用户手动设置ID。
env.socketTextStream("node5",9999).uid("socket-source")
.flatMap((String line,Collector<Tuple2<String,Integer>> out) -> {
String[] words = line.split(",");
for (String word : words) {
out.collect(new Tuple2<>(word, 1));
}
}).returns(Types.TUPLE(Types.STRING,Types.INT)).uid("flatmap")
.keyBy(tp -> tp.f0)
.sum(1).uid("sum")
.print().uid("print");
一般Flink程序进行升级时都涉及使用savepoint保存状态,然后升级后再基于savepoint进行状态恢复。下面以读取Socket中的数据进行WordCount为例来演示savepoint和算子uid的配置及使用,可以按照如下步骤进行测试。
编写代码并打包
java 代码
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); //读取socket数据做wordcount SingleOutputStreamOperator<String> lines = env.socketTextStream("node5", 9999).uid("socket-source"); SingleOutputStreamOperator<Tuple2<String, Integer>> tuple2 = lines.flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() { @Override public void flatMap(String s, Collector<Tuple2<String, Integer>> collector) throws Exception { String[] words = s.split(","); for (String word : words) { collector.collect(new Tuple2<>(word, 1)); } } }).uid("flatmap"); SingleOutputStreamOperator<Tuple2<String, Integer>> result = tuple2.keyBy(new KeySelector<Tuple2<String, Integer>, String>() { @Override public String getKey(Tuple2<String, Integer> stringIntegerTuple2) throws Exception { return stringIntegerTuple2.f0; } }).sum(1).uid("sum"); result.print().uid("print"); env.execute();scala 代码
val env = StreamExecutionEnvironment.getExecutionEnvironment //导入隐式转换 import org.apache.flink.streaming.api.scala._ // 读取socket数据做wordcount env.socketTextStream("node5", 9999).uid("socket-source") .flatMap(_.split(" ")).uid("flatMap") .map((_, 1)).uid("map") .keyBy(_._1) .sum(1).uid("sum") .print().uid("print") env.execute()
以上代码编写完成后,打包并上传到node5节点的/root/flink-jar-test目录中。
编写代码并打包
在node5节点配置flink-conf.yaml文件“Fault tolerance and checkpointing”部分配置savepoint保存的路径,如下。
state.savepoints.dir: hdfs://mycluster/flink-savepoints启动HDFS集群
#启动Zookeeper集群 [root@node3 ~]# zkServer.sh start [root@node4 ~]# zkServer.sh start [root@node5 ~]# zkServer.sh start #启动HDFS集群和Yarn集群 [root@node1 ~]# start-all.sh注意:启动HDFS集群后,如果有hdfs://mycluster/flink-savepoints目录,为了后续方便看出任务状态目录,最好删除该目录。
提交任务并统计状态结果
node5节点启动socket服务并在node5节点上提交Flink任务,提交Flink任务后,向Socket中输入数据,观察Flink任务webui中统计的结果。
#node5节点启动socket服务 [root@node5 ~]# nc -lk 9999 #node5节点提交Flink任务 [root@node5 conf]# cd /software/flink-1.17.1/bin/ [root@node5 bin]# ./flink run-application -t yarn-application -c com.wubaibao.flinkjava.code.chapter7.savepoints.SavePointTest /root/flink-jar-test/FlinkJavaCode-1.0-SNAPSHOT-jar-with-dependencies.jar提交任务后,向socket中输入以下数据:
hello,flink hello,flink hello,savepoint输入数据后,观察Flink WebUI统计的结果如下:

触发savepoint
执行如下命令执行savepoint 操作,将当前Flink程序的状态保存到对应路径中。
[root@node5 bin]# ./flink savepoint d2d863913b663ad70f60ee50d05aaa32 -yid application_1687674094408_0001以上命令注意以下几点:
- savepoint操作命令格式为:./flink savepoint []
- 如果savepoint path(target directory)在当前提交任务节点的flink-conf.yaml中配置了,就不需要再写上。
- 如果是基于Yarn中运行的Flink任务,Flink JobID 通过Flink WebUI查看,并且执行savepoint命令最后需要通过-yid参数指定yarn application ID 链接到Yarn Application中。
以上savepoint 命令执行后,可以看到HDFS中生成对应的savepoint路径:
此时,我们可以手动取消Flink任务或者通过命令停止Flink任务,操作如下:

#命令方式取消Flink任务 [root@node5 bin]# ./flink cancel d2d863913b663ad70f60ee50d05aaa32 -yid application_1687674094408_0001Flink任务启动后,可以继续向Socket中输入如下数据:
hello,flink hello,savepoint可以观察Flink WebUI对应统计的状态结果,可以看到Flink 任务已经成功从savepoint中恢复过来。

设置合适watermark
在Flink中,watermark是一种衡量事件时间进展的机制,watermark是一种特殊的数据记录,watermark本质就是一个时间戳,基于Flink接收到的事件时间(Event Time)计算得到,并且该时间标记会随着数据流往后流动,当Flink算子接收到Watermark(t)事件时,可以认为早于或等于t时刻的事件时间已经完全到达。
基于事件时间处理数据时定时器、流的关联、窗口触发都与watermark相关联,关于watermark的设置需要根据具体的业务数据延迟程度来决定,watermark延迟时间设置过大可能会导致内存使用过大或者窗口长时间不触发,这种情况下可以适当调小watermark大小,此外如果某个并行度中长时间没有数据到达,可以设置“WatermarkStrategy.xxx..withIdleness(Duration.ofSeconds(5))”来指定等待空闲时间自动推进watermark。