0%

Flink-checkpoint与高可用

Flink Checkpoint与高可用

Flink Checkpoint 受 Chandy-Lamport 分布式快照启发,可以保证数据的高可用。但是有些情况下,不见得一定有效:

Flink On Yarn 模式,某个 Container 发生 OOM 异常,这种情况程序直接变成失败状态,此时 Flink 程序虽然开启 Checkpoint 也无法恢复,因为程序已经变成失败状态,所以此时可以借助外部参与启动程序,比如外部程序检测到实时任务失败时,重新对实时任务进行拉起。

1.1. 2PC

1.1.1. Exactly-once VS At-least-once

算子做快照时,如果等所有输入端的barrier都到了才开始做快照,可保证算子的exactly-once;
如果为了降低延时而跳过对齐,从而继续处理数据,那么等barrier都到齐后做快照就是at-least-once了,因为这次的快照掺杂了下一次快照的数据,当作业失败恢复的时候,这些数据会重复作用系统,就好像这些数据被消费了两遍。

注:对齐只会发生在算子的上端是join操作以及上游存在partition或者shuffle的情况,对于直连操作类似map、flatMap、filter等还是会保证exactly-once的语义。

1.1.2. 端到端的Exactly once实现

2PC分为几个阶段: 开始事务->预提交->提交(或回滚)

为了保证Exactly once, Source和Sink必须支持Flink的2PC

当状态涉及到外部系统时,需要外部系统支持事务操作来配合Flink实现2PC协议,从而保证数据的exatly-once。
这个时候,sink算子除了将自己的state写到状态后端,还必须准备好事务提交。

  • 一旦所有的算子完成了它们的pre-commit,它们会要求一个commit。
  • 如果存在一个算子pre-commit失败了,本次事务失败,我们回滚到上次的checkpoint。
  • 一旦master做出了commit的决定,那么这个commit必须得到执行,就算宕机恢复也有继续执行。

1.1.2.1. pre-commit

pre-commit阶段起始于一次快照的开始,即master节点将checkpoint的barrier注入source端,barrier随着数据向下流动直到sink端。barrier每到一个算子,都会出发算子做本地快照。

precommit

当所有的算子都做完了本地快照并且回复master节点时,pre-commit阶段才算结束。这个时候,checkpoint已经成功,并且包含了外部系统的状态。如果作业失败,可以进行恢复。

precommit-success

1.1.2.2. commit

通知所有的算子这次checkpoint成功了,即2PC的commit阶段。source节点和window节点没有外部状态,所以这时它们不需要做任何操作。
而对于sink节点,需要commit这次事务,将数据写到外部系统。

commit

1.1.2.3. rollback

一旦任何一个算子的快照保存失败,则触发回滚,同样的sink算子也需要取消写入外部的数据

1.1.3. TwoPhaseCommitSinkFunction

为了简化2PC的实现成本,flink抽象了TwoPhaseCommitSinkFunction

  • beginTransaction。开始一次事务,在目的文件系统创建一个临时文件。接下来我们就可以将数据写到这个文件。
  • preCommit。在这个阶段,将文件flush掉,同时重起一个文件写入,作为下一次事务的开始。
  • commit。这个阶段,将文件写到真正的目的目录。值得注意的是,这会增加数据可视的延时。
  • abort。如果回滚,那么删除临时文件。

如果pre-commit成功了,但是commit没有到达算子旧宕机了,flink会将算子恢复到pre-commit时的状态,然后继续commit。

我们需要做的还有就是保证commit的幂等性,这可以通过检查临时文件是否还在来实现。

1.2. checkpoint

保留策略:

  • DELETE_ON_CANCELLATION 表示当程序取消时,删除 Checkpoint 存储文件。
  • RETAIN_ON_CANCELLATION 表示当程序取消时,保存之前的 Checkpoint 存储文件

默认情况下,Flink不会触发一次 Checkpoint 当系统有其他 Checkpoint 在进行时,也就是说 Checkpoint 默认的并发为1。

CheckpointCoordinator :

针对 Flink DataStream 任务,程序需要经历从 StreamGraph -> JobGraph -> ExecutionGraph -> 物理执行图四个步骤,其中在 ExecutionGraph 构建时,会初始化 CheckpointCoordinator。ExecutionGraph通过ExecutionGraphBuilder.buildGraph方法构建,在构建完时,会调用 ExecutionGraph 的enableCheckpointing方法创建CheckpointCoordinator

Flink Checkpoint 参数配置及建议:

  • 当 Checkpoint 时间比设置的 Checkpoint 间隔时间要长时,可以设置 Checkpoint 间最小时间间隔 。这样在上次 Checkpoint 完成时,不会立马进行下一次 Checkpoint,而是会等待一个最小时间间隔,然后在进行该次 Checkpoint。否则,每次 Checkpoint 完成时,就会立马开始下一次 Checkpoint,系统会有很多资源消耗 Checkpoint。
  • 如果Flink状态很大,在进行恢复时,需要从远程存储读取状态恢复,此时可能导致任务恢复很慢,可以设置 Flink Task 本地状态恢复。任务状态本地恢复默认没有开启,可以设置参数state.backend.local-recovery值为true进行激活。
  • Checkpoint保存数,Checkpoint 保存数默认是1,也就是保存最新的 Checkpoint 文件,当进行状态恢复时,如果最新的Checkpoint文件不可用时(比如HDFS文件所有副本都损坏或者其他原因),那么状态恢复就会失败,如果设置 Checkpoint 保存数2,即使最新的Checkpoint恢复失败,那么Flink 会回滚到之前那一次Checkpoint进行恢复。考虑到这种情况,用户可以增加 Checkpoint 保存数。
  • 建议设置的 Checkpoint 的间隔时间最好大于 Checkpoint 的完成时间。

下图是不设置 Checkpoint 最小时间间隔示例图,可以看到,系统一致在进行 Checkpoint,可能对运行的任务产生一定影响:

1.3. savepoint

注意:
使用DataStream进行开发,建议为每个算子定义一个 uid,这样我们在修改作业时,即使导致程序拓扑图改变,由于相关算子 uid 没有变,那么这些算子还能够继续使用之前的状态,如果用户没有定义 uid , Flink 会为每个算子自动生成 uid,如果用户修改了程序,可能导致之前的状态程序不能再进行复用。

Flink 在触发Savepoint 或者 Checkpoint时,会根据这次触发的类型计算出在HDFS上面的目录:

如果类型是 Savepoint,那么 其 HDFS 上面的目录为:Savepoint 根目录+savepoint-jobid前六位+随机数字,具体如下格式:

Checkpoint 目录为 chk-checkpoint ID,具体格式如下:

  • 使用 flink cancel -s 命令取消作业同时触发 Savepoint 时,会有一个问题,可能存在触发 Savepoint 失败。比如实时程序处于异常状态(比如 Checkpoint失败),而此时你停止作业,同时触发 Savepoint,这次 Savepoint 就会失败,这种情况会导致,在实时平台上面看到任务已经停止,但是实际实时作业在 Yarn 还在运行。针对这种情况,需要捕获触发 Savepoint 失败的异常,当抛出异常时,可以直接在 Yarn 上面 Kill 掉该任务。
  • 使用 DataStream 程序开发时,最好为每个算子分配 uid,这样即使作业拓扑图变了,相关算子还是能够从之前的状态进行恢复,默认情况下,Flink 会为每个算子分配 uid,这种情况下,当你改变了程序的某些逻辑时,可能导致算子的 uid 发生改变,那么之前的状态数据,就不能进行复用,程序在启动的时候,就会报错。
  • 由于 Savepoint 是程序的全局状态,对于某些状态很大的实时任务,当我们触发 Savepoint,可能会对运行着的实时任务产生影响,个人建议如果对于状态过大的实时任务,触发 Savepoint 的时间,不要太过频繁。根据状态的大小,适当的设置触发时间。
  • 当我们从 Savepoint 进行恢复时,需要检查这次 Savepoint 目录文件是否可用。可能存在你上次触发 Savepoint 没有成功,导致 HDFS 目录上面 Savepoint 文件不可用或者缺少数据文件等,这种情况下,如果在指定损坏的 Savepoint 的状态目录进行状态恢复,任务会启动不起来。

1.4. snapshot保存到哪里? 应该需要汇总到jobManager?

1.5. state backend

FsStateBackend

构造方法:
FsStateBackend(URI checkpointDataUri,boolean asynchronousSnapshots)

1 基于文件系统的状态管理器
2 如果使用,默认是异步
3 比较稳定,3个副本,比较安全。不会出现任务无法恢复等问题
4 状态大小受磁盘容量限制

存储方式:

  • State: TaskManager内存
  • checkpoint: 外部文件系统(本地或HDFS)

容量限制:

  • 单TaskManager上State总量不超过它的内存
  • 总大小不超过配置的文件系统容量

推荐使用场景:

  • 常规使用状态的作业,例如分钟级窗口聚合、join、窗口比较长、kv状态大;需要开启HA的作业
  • 可以用于生产场景

RocksDBStateBackend

状态数据先写入RocksDB,然后异步的将状态数据写入文件系统。正在进行计算的热数据存储在RocksDB,长时间才更新的数据写入磁盘中(文件系统)存储,体量比较小的元数据状态写入JobManager内存中(将工作state保存在RocksDB中,并且默认将checkpoint数据存在文件系统中)

目前唯一支持incremental的checkpoints的策略

构造方法:
RocksDBStateBackend(URI checkpointDataUri,boolean enableIncrementalCheckpointing)

存储方式:

  • State: TaskManager上的KV数据库(实际使用内存+硬盘)
  • Checkpoint: 外部文件系统(本地或HDFS)

容量限制:

  • 单TaskManager上的State总量不超过他的内存+磁盘
  • 单key最大2G
  • 总大小不超过配置的文件系统容量

推荐使用的场景:

  • 超大状态的作业,例如天级别窗口聚合;需要开启HA的作业;对状态读写性能要求不高的作业
  • 可以在生产环境使用

MemoryStateBackend

构造方法:
MemoryStateBackend(int maxStateSize, boolean asynchronousSnapshots)

主机内存中的数据可能会丢失,任务可能无法恢复

存储方式:

  • State: TaskManager内存
  • Checkpoint: JobManager内存

容量限制

  • 单个State maxStateSize默认5M
  • maxStateSize <= akka.frameSize 默认10M
  • 总大小不超过JobManager的内存

推荐使用场景:

  • 本地测试;几乎无状态的作业,比如ETL;JobManager不容易挂,或挂掉影响不大的情况
  • 不推荐在生产环境使用

1.6. checkpoint 与 savepoint

Checkpoint指定触发生成时间间隔后,每当需要触发Checkpoint时,会向Flink程序运行时的多个分布式的Stream Source中插入一个Barrier标记,这些Barrier会根据Stream中的数据记录一起流向下游的各个Operator。
当一个Operator接收到一个Barrier时,它会暂停处理Steam中新接收到的数据记录。
因为一个Operator可能存在多个输入的Stream,而每个Stream中都会存在对应的Barrier,该Operator要等到所有的输入Stream中的Barrier都到达。(对齐)
当所有Stream中的Barrier都已经到达该Operator,这时所有的Barrier在时间上看来是同一个时刻点(表示已经对齐),在等待所有Barrier到达的过程中,
Operator的Buffer中可能已经缓存了一些比Barrier早到达Operator的数据记录(Outgoing Records),这时该Operator会将数据记录(Outgoing Records)发射(Emit)出去,作为下游Operator的输入,
最后将Barrier对应Snapshot发射(Emit)出去作为此次Checkpoint的结果数据。

Checkpoint 是增量做的,每次的时间较短,数据量较小,只要在程序里面启用后会自动触发,用户无须感知;Checkpoint 是作业 failover 的时候自动使用,不需要用户指定。

Savepoint 是全量做的,每次的时间较长,数据量较大,需要用户主动去触发。Savepoint 一般用于程序的版本更新(详见文档),Bug 修复,A/B Test 等场景,需要用户指定。

保存的内容

  • 首先,Savepoint 包含了一个目录,其中包含(通常很大的)二进制文件,这些文件表示了整个流应用在 Checkpoint/Savepoint 时的状态。
  • 以及一个(相对较小的)元数据文件,包含了指向 Savapoint 各个文件的指针,并存储在所选的分布式文件系统或数据存储中。

目标

Savepoint 和 Checkpoint 的不同之处很像传统数据库中备份与恢复日志之间的区别。Checkpoint 的主要目标是充当 Flink 中的恢复机制,确保能从潜在的故障中恢复。相反,Savepoint 的主要目标是充当手动备份、恢复暂停作业的方法。

实现

Checkpoint 被设计成轻量和快速的机制。它们可能(但不一定必须)利用底层状态后端的不同功能尽可能快速地恢复数据。例如,基于 RocksDB 状态后端的增量检查点,能够加速 RocksDB 的 checkpoint 过程,这使得 checkpoint 机制变得更加轻量。相反,Savepoint 旨在更多地关注数据的可移植性,并支持对作业做任何更改而状态能保持兼容,这使得生成和恢复的成本更高

状态文件保留策略

Checkpoint默认程序删除,可以设置CheckpointConfig中的参数进行保留 。Savepoint会一直保存,除非用户删除 。

应用

  • 部署流应用的一个新版本,包括新功能、BUG 修复、或者一个更好的机器学习模型
  • 引入 A/B 测试,使用相同的源数据测试程序的不同版本,从同一时间点开始测试而不牺牲先前的状态
  • 在需要更多资源时扩容应用程序
  • 迁移流应用程序到 Flink 的新版本上,或者迁移到另一个集群

Flink数据一致性

一、综述

flink 通过内部依赖checkpoint 并且可以通过设置其参数exactly-once 实现其内部的一致性。但要实现其端到端的一致性,还必须保证
1、source 外部数据源可重设数据的读取位置
2、sink端 需要保证数据从故障恢复时,数据不会重复写入外部系统(或者可以逻辑实现写入多次,但只有一次生效的数据sink端)

二、sink 端到端实现方式

幂等操作:
一个操作,可以重复执行多次,但只导致一次结果更改,豁免重复操作执行就不起作用了,他的瑕疵 (在系统恢复的过程中,如果这段时间内多个更新或者插入导致状态不一致,但当数据追上就可以了)
(逻辑与、逻辑或等)具体理解参照自己以前写的文章。
事务写入:
事务应该具有四个属性:原子性、一致性、隔离性、持久性等。其具体的实现方式有两种
(1)、预写日志
简单易于实现,由于数据提前在状态后端中做了缓存,所以无论什么sink系统,都能用这种方式一批搞定,DataStream API提供了一个模板类:GenericWriteAheadSink,来实现这种事务性sink;
缺点:
1)、sink系统没说他支持事务。有可能出现一部分写入集群了。一部分没有写进去(如果实表,再写一次就写重复了)
2)、checkpoint做完了sink才去真正的写入(但其实得等sink都写完checkpoint才能生效,所以WAL这个机制jobmanager确定它写完还不算真正写完,还得有一个外部系统已经确认 完成的checkpoint)
2)、两阶段提交。 flink 真正实现exactle-once
对于每个checkpoint,sink 任务会启动一个事务,并将接下来所有接收的数据添加到事务中,然后将这些数据写入外部sink系统,但不提交他们(这里是预提交)。当checkpoint完成时的通知,它才正式提交事务,实现结果的真正写入;这种方式真正实现了exactly-once,它需要一个提供事务支持的外部sink系统,Flink提供了其具体实现(TwoPhaseCommitSinkFunction接口)

三、 2pc 对外部 sink的要求

1、外部sink系统必须事务支持,或者sink任务必须能够模拟外部系统上的事务;
2、在checkpoint的间隔期间里,必须能够开启一个事务,并接受数据写入。
3、在收到checkpoint完成通知之前,事务必须是“等待提交”的状态,在故障恢复的情况线,这可能需要一些时间。如果个时候sink系统关闭事务(例如超时了),那么未提交的数据就会丢失;
4、四年任务必选能够在进程失败后恢复事务
5、提交事务必须是幂等操作;

四、综上不同Source和sink的一致性保证:

在这里插入图片描述

flink 和kafka 端到端一致性(kafka(source+flink+kafka(sink)))
1、内部 – 利用checkpoint机制,把状态存盘,发生故障的时候可以恢复,保证内部的状态一致性
2、source – kafka consumer作为source,可以将偏移量保存下来,如果后续任务出现了故障,恢复的时候可以由连接器重置偏移量,重新消费数据,保证一致性;

1
2
3
4
kafka 0.8 和kafka 0.11 之后 通过以下配置将偏移量保存,恢复时候重新消费
kafka.setStartFromLatest();
kafka.setCommitOffsetsOnCheckpoints(false);
kafka 0.9 和kafka0.10 未验证是否支持这两个参数(todo)

3、sink FlinkkafkaProducer作为Sink,采用两阶段提交的sink,由下图可以看出flink 0.11 已经默认继承了TwoPhaseCommitSinkFunction
在这里插入图片描述
但我们需要在参数种传入指定语义,它默认时还是at-least-once
此外我们还需要进行一些producer的容错配置:
(1)除了启用Flink的检查点之外,还可以通过将适当的semantic参数传递给FlinkKafkaProducer011(FlinkKafkaProducer对于Kafka> = 1.0.0版本)
(2)来选择三种不同的操作模式
1)、Semantic.NONE 代表at-mostly-once语义
2)、Semantic.AT_LEAST_ONCE(Flink默认设置
3)、Semantic.EXACTLY_ONCE 使用Kafka事务提供一次精确的语义,每当您使用事务写入Kafka时
(3)、请不要忘记消费kafka记录任何应用程序设置所需的设置isolation.leva(read_committed 或者read_uncommitted-后者是默认)
read_committed,只是读取已经提交的数据。

应用;
Semantic.EXACTLY_ONCE依赖与下游系统能支持事务操作.以0.11kafka为例.
transaction.max.timeout.ms 最大超市时长,默认15分钟,如果需要用exactly语义,需要增加这个值。(因为它小于transaction.timeout.ms )
isolation.level 如果需要用到exactly语义,需要在下级consumerConfig中设置read-commited [read-uncommited(默认值)]
transaction.timeout.ms 默认为1hour

其参数对应关系为 和一些报错问题
checkpoint间隔<transaction.timeout.ms<transaction.max.timeout.ms

参考:https://www.cnblogs.com/createweb/p/11971846.html

注意:
1、Semantic.EXACTLY_ONCE 模式每个FlinkKafkaProducer011实例使用一个固定大小的KafkaProducers池。每个检查点使用这些生产者中的每一个。如果并发检查点的数量超过池大小,FlinkKafkaProducer011 将引发异常,并使整个应用程序失败。请相应地配置最大池大小和最大并发检查点数。

2、Semantic.EXACTLY_ONCE采取所有可能的措施,不要留下任何挥之不去的数据,否则这将有碍于消费者更多地阅读Kafka主题。但是,如果flink应用程序在第一个检查点之前失败,则在重新启动此类应用程序后,系统种将没有有关先前池大小信息,因此,在第一个检查点完成前按比例缩小Flink应用程序的FlinkKafkaProducer011.SAFE_SCALE_DOWN_FACTOR

1
2
3
4
5
6
7
8
9
10
//1。设置最大允许的并行checkpoint数,防止超过producer池的个数发生异常
env.getCheckpointConfig.setMaxConcurrentCheckpoints(5)
//2。设置producer的ack传输配置
// 设置超市时长,默认15分钟,建议1个小时以上
producerConfig.put(ProducerConfig.ACKS_CONFIG, 1)
producerConfig.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, 15000)

//3。在下一个kafka consumer的配置文件,或者代码中设置ISOLATION_LEVEL_CONFIG-read-commited
//Note:必须在下一个consumer中指定,当前指定是没用用的
kafkaonfigs.setProperty(ConsumerConfig.ISOLATION_LEVEL_CONFIG,"read_commited")

完整应用代码:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
package com.shufang.flink.connectors

import java.util.Properties
import org.apache.flink.api.common.serialization.SimpleStringSchema
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.streaming.api.functions.timestamps.BoundedOutOfOrdernessTimestampExtractor
import org.apache.flink.streaming.api.scala._
import org.apache.flink.streaming.api.windowing.time.Time
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer.Semantic
import org.apache.flink.streaming.connectors.kafka._
import org.apache.flink.streaming.util.serialization.KeyedSerializationSchemaWrapper
import org.apache.kafka.clients.consumer.ConsumerConfig
import org.apache.kafka.clients.producer.ProducerConfig
import org.apache.kafka.common.serialization.StringDeserializer

object KafkaSource01 {
def main(args: Array[String]): Unit = {
val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment

env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)

//这是checkpoint的超时时间
//env.getCheckpointConfig.setCheckpointTimeout()
//设置最大并行的chekpoint
env.getCheckpointConfig.setMaxConcurrentCheckpoints(5)
env.getCheckpointConfig.setCheckpointInterval(1000) //增加checkpoint的中间时长,保证可靠性


/**
* 为了保证数据的一致性,我们开启Flink的checkpoint一致性检查点机制,保证容错
*/
env.enableCheckpointing(60000)

/**
* 从kafka获取数据,一定要记得添加checkpoint,能保证offset的状态可以重置,从数据源保证数据的一致性
* 保证kafka代理的offset与checkpoint备份中保持状态一致
*/

val kafkaonfigs = new Properties()

//指定kafka的启动集群
kafkaonfigs.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")
//指定消费者组
kafkaonfigs.setProperty(ConsumerConfig.GROUP_ID_CONFIG, "flinkConsumer")
//指定key的反序列化类型
kafkaonfigs.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer].getName)
//指定value的反序列化类型
kafkaonfigs.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, classOf[StringDeserializer].getName)
//指定自动消费offset的起点配置
// kafkaonfigs.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "latest")


/**
* 自定义kafkaConsumer,同时可以指定从哪里开始消费
* 开启了Flink的检查点之后,我们还要开启kafka-offset的检查点,通过kafkaConsumer.setCommitOffsetsOnCheckpoints(true)开启,
* 一旦这个检查点开启,那么之前配置的 auto-commit-enable = true的配置就会自动失效
*/
val kafkaConsumer = new FlinkKafkaConsumer[String](
"console-topic",
new SimpleStringSchema(), // 这个schema是将kafka的数据应设成Flink中的String类型
kafkaonfigs
)

// 开启kafka-offset检查点状态保存机制
kafkaConsumer.setCommitOffsetsOnCheckpoints(true)

// kafkaConsumer.setStartFromEarliest()//
// kafkaConsumer.setStartFromTimestamp(1010003794)
// kafkaConsumer.setStartFromLatest()
// kafkaConsumer.setStartFromSpecificOffsets(Map[KafkaTopicPartition,Long]()

// 添加source数据源
val kafkaStream: DataStream[String] = env.addSource(kafkaConsumer)

kafkaStream.print()

val sinkStream: DataStream[String] = kafkaStream.assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor[String](Time.seconds(5)) {
override def extractTimestamp(element: String): Long = {
element.split(",")(1).toLong
}
})


/**
* 通过FlinkkafkaProduccer API将stream的数据写入到kafka的'sink-topic'中
*/
// val brokerList = "localhost:9092"
val topic = "sink-topic"
val producerConfig = new Properties()
producerConfig.put(ProducerConfig.ACKS_CONFIG, new Integer(1)) // 设置producer的ack传输配置
producerConfig.put(ProducerConfig.TRANSACTION_TIMEOUT_CONFIG, Time.hours(2)) //设置超市时长,默认1小时,建议1个小时以上

/**
* 自定义producer,可以通过不同的构造器创建
*/
val producer: FlinkKafkaProducer[String] = new FlinkKafkaProducer[String](
topic,
new KeyedSerializationSchemaWrapper[String](SimpleStringSchema),
producerConfig,
Semantic.EXACTLY_ONCE
)

// FlinkKafkaProducer.SAFE_SCALE_DOWN_FACTOR
/** *****************************************************************************************************************
* * 出了要开启flink的checkpoint功能,同时还要设置相关配置功能。
* * 因在0.9或者0.10,默认的FlinkKafkaProducer只能保证at-least-once语义,假如需要满足at-least-once语义,我们还需要设置
* * setLogFailuresOnly(boolean) 默认false
* * setFlushOnCheckpoint(boolean) 默认true
* * come from 官网 below:
* * Besides enabling Flink’s checkpointing,you should also configure the setter methods setLogFailuresOnly(boolean)
* * and setFlushOnCheckpoint(boolean) appropriately.
* ******************************************************************************************************************/

producer.setLogFailuresOnly(false) //默认是false


/**
* 除了启用Flink的检查点之外,还可以通过将适当的semantic参数传递给FlinkKafkaProducer011(FlinkKafkaProducer对于Kafka> = 1.0.0版本)
* 来选择三种不同的操作模式:
* Semantic.NONE 代表at-mostly-once语义
* Semantic.AT_LEAST_ONCE(Flink默认设置)
* Semantic.EXACTLY_ONCE:使用Kafka事务提供一次精确的语义,每当您使用事务写入Kafka时,
* 请不要忘记为使用Kafka记录的任何应用程序设置所需的设置isolation.level(read_committed 或read_uncommitted-后者是默认值)
*/

sinkStream.addSink(producer)

env.execute("kafka source & sink")
}
}

[参考文献]

  1. Flink Checkpoint、Savepoint配置与实践
  2. Flink 小贴士 (2):Flink 如何管理 Kafka 消费位点
  3. Flink实时计算-深入理解Checkpoint和Savepoint
  4. Lightweight Asynchronous Snapshots for Distributed Dataflows: 分布式数据流轻量级异步快照