Flink 手写AT_LEAST_ONCE语义的Source
需求
老王交给个任务,希望我用Flink实现一个简单的数据ETL。需要实时拉取数据,保证数据不丢。我想了想整理这篇文章。仅供提供思路,并把相关知识串起来。
Flink本身提供了CDC功能更加强大方便。虽然如此,但是我想通过 SourceFunction 实现Mysql实时数据拉取。简单思路:
为了保证最少一次读取数据(数据不丢),使用状态存储查询到最新截止时间。这里需要自定义实现数据库连接。使用官方jdbc-sink写入mysql
使用到知识点
自定义 SourceFunction 状态 Operator State (non-keyed state) 检查点 CheckPoint 状态后端 StateBackend 数据库写入 Mysql Sink
1. 自定义 SourceFunction
Flink 对数据的产生通过提供了统一的接口,方便用户使用。
首先看一张类继承依赖图:
我们只需要重点关注下面3个类就好,首先看下他们的接口定义:
SourceFunction:Flink中所有流数据源的基本接口 RichSourceFunction :用于实现可以访问上下文信息的数据源的基类 RichParallelSourceFunction:用于实现并行数据源的基类
1.1 SourceFunction接口
再来看看SourceFunction接口的常用方法:
@Public
public interface SourceFunction<T> extends Function, Serializable {
/**
* 启动该数据源发送元素
*/
void run(SourceContext<T> ctx) throws Exception;
/**
* 取消这个源的数据发送
*/
void cancel();
}
SourceFunction实现例子:
/**
* 不带上下文 不支持并发
*/
public class E02CustomizeSource implements SourceFunction<OrderInfo> {
//运行标志位
private volatile boolean isRunning = true;
@Override
public void run(SourceContext<OrderInfo> ctx) throws Exception {
while (isRunning) {
OrderInfo orderInfo = new OrderInfo();
orderInfo.setOrderNo(System.currentTimeMillis() + "" + new Random().nextInt(10));
orderInfo.setItemId(1000);
orderInfo.setItemName("iphone 12pro max");
orderInfo.setPrice(9800D);
orderInfo.setQty(1);
ctx.collect(orderInfo);
Thread.sleep(new Random().nextInt(10) * 1000);
}
}
@Override
public void cancel() {
isRunning = false;
}
}
从例子可以看出 run
方法中只要 isRunning!=false
就可以一直运行下 , cancel
可以修改这个标志位,
SourceContext ctx ctx.collect(元素); 我们通过该方法发送元素数据
1.2 RichSourceFunction 抽象方法
@Public
public abstract class RichSourceFunction<OUT> extends AbstractRichFunction implements SourceFunction<OUT> {
private static final long serialVersionUID = 1L;
}
该方法核心方法继承自 SourceFunction
但是为了提供上下文信息又继承了 AbstractRichFunction
且额外提供了 资源初始化:open()
和 资源清理 close()
分别在 run()
运行前执行和运行后运行,常用比如Mysql连接打开和关闭。
@Public
public abstract class AbstractRichFunction implements RichFunction, Serializable {
private transient RuntimeContext runtimeContext;
//该方法可以获取上下文
@Override
public RuntimeContext getRuntimeContext() {
if (this.runtimeContext != null) {
return this.runtimeContext;
} else {
throw new IllegalStateException("The runtime context has not been initialized.");
}
}
//资源初始化
@Override
public void open(Configuration parameters) throws Exception {}
//资源关闭
@Override
public void close() throws Exception {}
}
RichSourceFunction 实现实例:
/**
* 带上下文 不支持并发
* 而且增加了两个资源打开关闭方法 可以再run之前调用
*/
public class E03CustomizeRichSource extends RichSourceFunction<OrderInfo> {
private transient Connection connection;
private transient PreparedStatement ps;
//运行标志位
private volatile boolean isRunning = true;
/**
* run 之前调用
*
* @param parameters
* @throws Exception
*/
@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
//通过上下问获取全局配置对象且读取配置信息
ParameterTool parameterTool = (ParameterTool) getRuntimeContext().getExecutionConfig().getGlobalJobParameters();
String url = parameterTool.get("url");
String username = parameterTool.get("username");
String password = parameterTool.get("password");
//初始化
Class.forName("com.mysql.cj.jdbc.Driver");
connection = DriverManager.getConnection(url, username, password);
ps = connection.prepareStatement("select * from order_info");
}
/**
* run 之后调用
*
* @throws Exception
*/
@Override
public void close() throws Exception {
super.close();
ps.close();
connection.close();
}
@Override
public void run(SourceContext<OrderInfo> ctx) throws Exception {
try (ResultSet resultSet = ps.executeQuery()) {
while (resultSet.next()) {
//todo ...
//ctx.
}
}
}
@Override
public void cancel() {
isRunning = false;
}
}
1.3 RichParallelSourceFunction 抽象方法
@Public
public abstract class RichParallelSourceFunction<OUT> extends AbstractRichFunction
implements ParallelSourceFunction<OUT> {
private static final long serialVersionUID = 1L;
}
从名字就可以知道这是一个并行数据源,他只是把 SourceFunction
替换了 ParallelSourceFunction
。SourceFunction
和 RichSourceFunction
并行度设置只能为 1 否则会报错**_ The parallelism of non parallel operator must be 1. __ 。_** RichParallelSourceFunction
则不受影响。
java.lang.IllegalArgumentException: The parallelism of non parallel operator must be 1.
at org.apache.flink.util.Preconditions.checkArgument(Preconditions.java:142)
at org.apache.flink.api.common.operators.util.OperatorValidationUtils.validateParallelism(OperatorValidationUtils.java:37)
at org.apache.flink.streaming.api.datastream.DataStreamSource.setParallelism(DataStreamSource.java:103)
at flinklearning._1source.DemoMain.main(DemoMain.java:30)
1.4 需求选型
我们知道我们查询数据库并行查询,会查询重复数据,所以没必须要使用 RichParallelSourceFunction
。我们需要使用mysql。需要进行资源操作,所以最好使用 RichSourceFunction
。
2. 检查点 CheckPoint & 状态 Operator State
Flink 是有状态的流处理框架。本人理解:状态是某个时间节点元素产生事件的描述。比如:你想管理历史数据的时候,你可以将该历史数据存储在状态里面,这样你可以在未来某个时间访问。

2.1 状态的持久性
Flink 通过 CheckPoint(检查点) 将 state(状态) 持久化在状态后端。如果程序失败或异常重启,可以通过持久化的检查点恢复到当时状态。常用后端:
MemoryStateBackend (默认) FsStateBackend RocksDBStateBackend
_ _ _
2.2 检查点触发
Flink 默认是不开启检查点,用户可以 StreamExecutionEnvironment.enableCheckpointing(2000``)
设置。Flink 会根据用户设置的时间 周期有序的从source往数据流中插入** barrier(栏栅) **
如:通过下图我们初步了解下 检查点机制 并行度为一情况
时间线 1 Source
端网数据流中插入checkPoint 1时间线 2 checkPoint 1
进入Map
触发State
持久化 ,并继续往下游算子数据流中传送。时间线3 checkPoint 1
进入Sink
触发State
持久化 ,整个检查点触发完毕。
2.3 Operator State 使用
Operator State
,只支持 ListState
状态存储结构和广播状态结构。其实还有另一种类型状态,这里我们先只记住这个,有机会后面单独讲讲区别和使用。使用Operator State
需要结合CheckpointedFunction
接口使用。算子实现该接口,重写其中的两个方法:
snapshotState 快照状态 触发检查点快照的时候 initializeState 初始化状态 再构造函数后调用
public interface CheckpointedFunction{
//触发检查点的时会调用此方法
void snapshotState(FunctionSnapshotContext context) throws Exception;
//创建函数的时候回嗲用此方法,可以获取状态信息 恢复嘻嘻
void initializeState(FunctionInitializationContext context) throws Exception;
}
代码示例:
public class E05ReadMysqlRichSource02 implements CheckpointedFunction {
/**
* 存储 state 的变量.
*/
private ListState<LocalDateTime> startDateTimeState;
/**
* 检查点快照恢复
*
* @param context
* @throws Exception
*/
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
System.err.println("快照时候执行 snapshotState---------------");
startDateTimeState.clear();
startDateTimeState.add(startDateTime);
System.out.println("快照结果: " + startDateTime);
}
/**
* 启动时执行快照恢复
*
* @param context
* @throws Exception
*/
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
System.err.println("第二执行 initializeState---------------");
//定义状态描述
ListStateDescriptor<LocalDateTime> startDateTimeStateDescriptor = new ListStateDescriptor<LocalDateTime>("startDateTimeState", TypeInformation.of(LocalDateTime.class));
startDateTimeState = context.getOperatorStateStore().getListState(startDateTimeStateDescriptor);
//从状态中恢复
if (context.isRestored()) {
for (LocalDateTime localDateTime : startDateTimeState.get()) {
startDateTime = localDateTime;
}
System.out.println("恢复结果:" + startDateTime);
}
}
}
2.4 需求分析
上面我们了解快照的作用,和怎么使用状态。我们就可以解决需求中 最新查询截止时间 不会因为程序故障丢失,实现最少一次读取。
3. 理解自定义实现的SourceFunction生命周期
我们尝试来编写 MysqlReadSourceFunction 首先继承 RichSourceFunction
实现数据发送;实现**CheckpointedFunction**
来保障数据完整性。但是提供了那么方法谁先执行谁后执行,关系是什么?
/**
* 实现自定实现 Mysql 数据读取器
* 支持 定点读取和切好一次读取
*/
public class MysqlReadSourceFunction extends RichSourceFunction<TraceSegmentRecordInfo> implements CheckpointedFunction {
public MysqlReadSourceFunction() {
super();
System.err.println("第一执行 1构造方法");
}
@Override
public void open(Configuration parameters) throws Exception {
System.err.println("第三执行 open---------------");
super.open(parameters);
}
@Override
public void close() throws Exception {
System.err.println("第五执行 close---------------");
super.close();
}
@Override
public void run(SourceContext<TraceSegmentRecordInfo> ctx) throws Exception {
System.err.println("第三四 run---------------");
}
@Override
public void cancel() {
System.err.println("取消时候执行 cancel---------------");
}
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
System.err.println("快照时候执行 snapshotState---------------");
}
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
System.err.println("第二执行 initializeState---------------");
}
}
首先肯定是构造器 其次会调用 initializeState 方法,进行数据状态恢复 再而 open() 方法 初始化资源 紧接着 就进入真正执行 执行数据发送 run() 核心方法 此时如果程序取消了,会调用 cancel() 可以在里面关闭 run() 里面的循环,它会在线程中断前执行。 接着会调用 清理 close() 最后结束
3. 需求实现,就是填写代码了
/**
* 实现自定实现 Mysql 数据读取器
* 支持 定点读取和至少一次
*/
public class E05ReadMysqlRichSource02 extends RichSourceFunction<TraceSegmentRecordInfo> implements CheckpointedFunction {
//常量名称
public static final String url = "url";
public static final String username = "username";
public static final String password = "password";
//数据库连接资源对象
private transient Connection connection;
private transient PreparedStatement statement;
//查询SQL
private String querySql;
//是否继续运行
private volatile boolean isRunning = true;
//查询开始时间
private volatile LocalDateTime startDateTime;
//时间增长步长 单位分钟
private volatile int timeStep = 1;
/**
* 存储 state 的变量.
*/
private ListState<LocalDateTime> startDateTimeState;
public E05ReadMysqlRichSource02(String querySql, LocalDateTime startDateTime) {
super();
//设置查询SQL
this.querySql = querySql;
//初始化查询其实时间
this.startDateTime = startDateTime;
}
/**
* 在run方法前执行
* 进行资源初始化
*
* @param parameters
* @throws Exception
*/
@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
//获取系统配置
ParameterTool parameterTool = (ParameterTool) getRuntimeContext().getExecutionConfig().getGlobalJobParameters();
//创建连接资源
Class.forName("com.mysql.cj.jdbc.Driver");
connection = DriverManager.getConnection(parameterTool.get(url), parameterTool.get(username), parameterTool.get(password));
statement = connection.prepareStatement(querySql);
}
/**
* 资源进行关闭 在线程中断前执行
*
* @throws Exception
*/
@Override
public void close() throws Exception {
super.close();
System.err.println("第五执行 close---------------");
isRunning = false;
//关闭资源连接
if (statement != null) {
statement.close();
}
if (connection != null) {
connection.close();
}
}
/**
* 在支持只支持不会并行执行
*
* @param ctx
* @throws Exception
*/
@Override
public void run(SourceContext<TraceSegmentRecordInfo> ctx) throws Exception {
final DateTimeFormatter dtf = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss");
//获取检查点锁
final Object lock = ctx.getCheckpointLock();
//循环遍历
while (isRunning) {
//按照 时间步长生成 查询起止时间
String startDateTimeStr = dtf.format(startDateTime);
LocalDateTime endDateTime = startDateTime.plusMinutes(timeStep);
if (endDateTime.isAfter(LocalDateTime.now())) {
//最新截止时间 大于系统时间休眠后跳出循环
Thread.sleep(timeStep * 60000);
continue;
}
String endDateTimeStr = dtf.format(endDateTime);
//设置参数
statement.setString(1, startDateTimeStr);
statement.setString(2, endDateTimeStr);
//查询数据且封装参数
List<TraceSegmentRecordInfo> containerList = new ArrayList<>();
ResultSet resultSet = statement.executeQuery();
while (resultSet.next()) {
TraceSegmentRecordInfo info = new TraceSegmentRecordInfo();
info.setSegmentId(resultSet.getString("segment_id"));
info.setTraceId(resultSet.getString("trace_id"));
info.setServiceName(resultSet.getString("service_name"));
info.setServiceIp(resultSet.getString("service_ip"));
info.setEndpointName(resultSet.getString("endpoint_name"));
info.setDataBinary(resultSet.getString("data_binary"));
info.setTimeBucket(resultSet.getLong("time_bucket"));
info.setStartTime(resultSet.getLong("start_time"));
info.setEndTime(resultSet.getLong("end_time"));
info.setLatency(resultSet.getInt("latency"));
info.setIsError(resultSet.getInt("is_error"));
info.setCreateDate(resultSet.getDate("create_time"));
info.setStatement(resultSet.getString("statement"));
containerList.add(info);
}
System.out.println(String.format("执行结果[%s - %s] 查询条数%s", startDateTimeStr, endDateTimeStr, containerList.size()));
if (CollectionUtils.isNotEmpty(containerList)) {
for (Iterator<TraceSegmentRecordInfo> iterator = containerList.iterator(); iterator.hasNext(); ) {
//发送数据
ctx.collect(iterator.next());
}
}
synchronized (lock) {
//推送完毕 最后使用endDateTime 设置为起始时间
startDateTime = endDateTime;
}
}
}
@Override
public void cancel() {
System.err.println("取消时候执行 cancel---------------");
//设置停止run方法循环
isRunning = false;
}
/**
* 检查点快照恢复
*
* @param context
* @throws Exception
*/
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
System.err.println("快照时候执行 snapshotState---------------");
//快照最新状态
startDateTimeState.clear();
startDateTimeState.add(startDateTime);
System.out.println("快照结果: " + startDateTime);
}
/**
* 启动时执行快照恢复
*
* @param context
* @throws Exception
*/
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
System.err.println("第二执行 initializeState---------------");
//定义状态描述
ListStateDescriptor<LocalDateTime> startDateTimeStateDescriptor = new ListStateDescriptor<LocalDateTime>("startDateTimeState", TypeInformation.of(LocalDateTime.class));
startDateTimeState = context.getOperatorStateStore().getListState(startDateTimeStateDescriptor);
//从状态中恢复,
if (context.isRestored()) {
//只有重状态恢复 才会进来,正常启动 isRestored = false
for (LocalDateTime localDateTime : startDateTimeState.get()) {
startDateTime = localDateTime;
}
System.out.println("恢复结果:" + startDateTime);
}
}
}
public class DemoMain {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
//自己实现带环境的系统参数解析
ParameterTool parameterTool = ParameterToolEnvironmentUtils.createParameterTool(args);
env.getConfig().setGlobalJobParameters(parameterTool);
//开启检查点
env.enableCheckpointing(2000L);
env.getCheckpointConfig().enableExternalizedCheckpoints(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
//设置文件系统作为状态后端
FsStateBackend fsStateBackend = new FsStateBackend("file:///usr/flink-1.12.1/tmp");
env.setStateBackend(fsStateBackend);
// 执行结果[2021-02-09 07:39:01 - 2021-02-09 07:40:01] 查询条数177
String sql = "SELECT * from trace_segment_record where create_time >= ? AND create_time < ?";
//source
DataStreamSource<TraceSegmentRecordInfo> source2 = env.addSource(new E05ReadMysqlRichSource02(sql, LocalDateTime.of(2021, 01, 8, 0, 0, 01))).setParallelism(2);
//过滤
SingleOutputStreamOperator<TraceSegmentRecordInfo> filter = source2.filter(new FilterFunction<TraceSegmentRecordInfo>() {
@Override
public boolean filter(TraceSegmentRecordInfo value) throws Exception {
//过滤等于null的数据
return value.getEndpointName() != null;
}
});
//sink
JdbcExecutionOptions executionOptions = JdbcExecutionOptions.builder().withBatchIntervalMs(2000L).withBatchSize(5000).build();
JdbcConnectionOptions connectionOptions = (new JdbcConnectionOptions.JdbcConnectionOptionsBuilder())
.withUrl("jdbc:mysql://xxxx:3306/yto_data_pipeline_init?useUnicode=true&characterEncoding=UTF-8&useSSL=false&serverTimezone=Asia/Shanghai")
.withDriverName("com.mysql.cj.jdbc.Driver")
.withUsername("xxx")
.withPassword("xxxx").build();
final String sinkSql = "REPLACE INTO trace_segment_record (segment_id, trace_id, service_name, service_ip, endpoint_name, start_time, end_time,latency,is_error, data_binary,time_bucket,create_time,statement) VALUE(?,?,?,?,?,?,?,?,?,?,?,?,?)";
filter.addSink(JdbcSink.sink(sinkSql, (ps, v) -> {
ps.setString(1, v.getSegmentId());
ps.setString(2, v.getTraceId());
ps.setString(3, v.getServiceName());
ps.setString(4, v.getServiceIp());
ps.setString(5, v.getEndpointName());
ps.setLong(6, v.getStartTime());
ps.setLong(7, v.getEndTime());
ps.setLong(8, v.getLatency());
ps.setInt(9, v.getIsError());
ps.setString(10, v.getDataBinary());
ps.setLong(11, v.getTimeBucket());
ps.setTimestamp(12, new Timestamp(v.getCreateDate().getTime()));
ps.setString(13, v.getStatement());
}, executionOptions, connectionOptions)).uid("sink_mysql");
env.execute("aaa");
}
}
3.1 中断写入演示
我在执行到 2021-01-08 00:32:01
取消了任务,发现系统已经存储快照了检查点。
我们使用检查点重启项目在观察日志是否从 2021-01-08 00:32:01
开始打印!
其实从日志可以看到 我打印了日志不代表那个时间节点已经持久化到数据,数据库,
这里取得上一次的时间节点 2021-01-08 00:30:01
数据也的确是从上次检查点开始查询的。
其实大家也发现数据重复读取了,这个也是文章开头最少一次读取,部分重复,最后是通过数据唯一性插入语句实现数据一致性的。





