暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Flink 手写AT_LEAST_ONCE语义的Source

OfNull 2021-03-03
809

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<Textends FunctionSerializable {
    /**
  * 启动该数据源发送元素
     */

 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
 可以修改这个标志位,

  • SourceContextctx
    • ctx.collect(元素);   我们通过该方法发送元素数据

1.2 RichSourceFunction 抽象方法

@Public
public abstract class RichSourceFunction<OUTextends AbstractRichFunction implements SourceFunction<OUT{

 private static final long serialVersionUID = 1L;
}

该方法核心方法继承自 SourceFunction
 但是为了提供上下文信息又继承了 AbstractRichFunction
 且额外提供了 资源初始化:open()
  和  资源清理 close()
  分别在 run()
 运行前执行和运行后运行,常用比如Mysql连接打开和关闭。

@Public
public abstract class AbstractRichFunction implements RichFunctionSerializable {
 
    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<OUTextends 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 是有状态的流处理框架。本人理解:状态是某个时间节点元素产生事件的描述。比如:你想管理历史数据的时候,你可以将该历史数据存储在状态里面,这样你可以在未来某个时间访问。

状态 (1).png

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<TraceSegmentRecordInfoimplements 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---------------");
    }
}


  1. 首先肯定是构造器
  2. 其次会调用 initializeState 方法,进行数据状态恢复
  3. 再而 open() 方法 初始化资源
  4. 紧接着 就进入真正执行 执行数据发送 run() 核心方法
  5. 此时如果程序取消了,会调用 cancel() 可以在里面关闭 run() 里面的循环,它会在线程中断前执行。
  6. 接着会调用 清理 close()
  7. 最后结束

3. 需求实现,就是填写代码了

/**
 * 实现自定实现 Mysql 数据读取器
 * 支持 定点读取和至少一次
 */

public class E05ReadMysqlRichSource02 extends RichSourceFunction<TraceSegmentRecordInfoimplements 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(20210180001))).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(12new 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
数据也的确是从上次检查点开始查询的。

其实大家也发现数据重复读取了,这个也是文章开头最少一次读取,部分重复,最后是通过数据唯一性插入语句实现数据一致性的。


文章转载自OfNull,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论