改版通知

巨人肩膀网站已全新改版。若您仍依赖旧站功能或数据,欢迎联系我们,我们会协助处理。联系我们

万字长文Flinkcdc源码精讲推荐收藏

ckckck2025年1月10日52 浏览

编者荐语
相对比较全面的一篇文章。

以下文章来源于857Hub,作者徐志文

857Hub 专注于大数据开发、数据架构之路,热衷于分享Hadoop、Flink、Spark、Doris、实时数仓、推荐等精品干货!

前言

flink-cdc源码地址: https://github.com/ververica/flink-cdc-connectors

flink-cdc不在flink项目中,在flink1.11之后flink引入cdc功能,下面我们以源码深入了解flink-cdc实现原理。

我们主要以flink-cdc-mysql为主,其余代码基本差不太多。

事先需要先简单了解一下debezium相关原理,flink-cdc是基于debezium实现的。

一点建议

  • 在阅读源码的时候,我们应该带着问题去思考,然后一步一步去阅读源码。在阅读源码的过程中,不要被一些不重要的点占用过多的时间精力,并且一遍两遍是不会让我们有一个清晰的印象的。毕竟别人多少年多少人的开发,看一两遍就可以理解的。在阅读某个框架源码之前,我们应该已经对该框架原理有一定的理解,然后根据我们的理解去验证他是代码实现的样子,或者带着思考去阅读,为什么这么实现,这么实现的好处是什么等。其实代码都是一样的,只不过是每个人的实现方式不同,考虑的问题不同而已。

  • 要有一定的Java基础,熟悉多线程,了解开发使用的相关接口(或者自己看了介绍之后很容易理解)。如果基础不牢,更多的是建议先从基础学习,然后写一写代码测试,比如多线程的时候怎么做交互等,自己写一写,在后面阅读源码的时候会更容易理解里面内容。

  • 该内容要首先对cdc有一定的了解,知道cdc的相关原理,flink-cdc的实现基于debezium实现,该框架是开源的,可以先去了解一下,这样对于我们后面内容会更容易理解。

谨记:阅读的时候抓住重点,不要被不重要的内容占用时间。

一. 项目结构(mysql-cdc为主)

1. 目录结构

  • 带有test项目都是用于测试的项目。

  • 后缀带有cdc的表示一个database的连接器,区分sql与api形式。

  • flink-format-changelog-json:用于解析json成RowData的模块。

  • flink-connector-debezium:该模块封装debezium以及相关核心代码实现,并且修改了debezium的部分源码。

  • 每个项目中都有test目录,里面有相关的测试代码,可以自行测试代码debug。

2. mysql项目源码包结构

  • debezium:debezium用到的相关类。

  • schema:mysql schema(表结构)相关代码。

  • source:mysql-cdc source实现代码,包括全量读mysql,分割器,读取器等相关。

  • table:cdc table实现代码主要以table dynamic factory的实现。

  • resources:该目录用于spi方式动态加载table factory,用于sql创建table找到对应的工厂类。

二. mysql-cdc源码 - SourceFunction的单并行度的实现

  • 基于RichSourceFunction的,单并行读取1.11之前的source接口,已被标记Deprecated。

  • 基于Source的,多并行度,1.11之后新出的source接口,实现要更复杂。

我们主要根据单并行度源码进行讲解,这样更方便理解。

具体入手我们可以根据文档中的创建source的类来一点一点走。

java 复制代码
MySqlSource通过构建者模式(23种设计模式)构建,我们只需要知道我们可以设置哪些参数即可,这个比较容易理解。

// 通过构建者方式配置任务启动时候所需要的参数
public static class Builder<T> { // MySqlSource内部类
    private int port = 3306; // default 3306 port
    private String hostname;
    private String[] databaseList;
    private String username;
    private String password;
    private Integer serverId;
    private String serverTimeZone; // 时区
    private String[] tableList;
    private Properties dbzProperties; // 传入的dbz引擎所需的属性
    private StartupOptions startupOptions = StartupOptions.initial(); // 用于控制开始binlog开始消费位置的参数
    private DebeziumDeserializationSchema<T> deserializer; // 用于对数据解析成什么样子如json,String等,定义序列化方式

    // 上面参数配置完成,通过build构建sourceFunction,主要将配置信息封装到properties中,这里面的参数主要是debezium所需要启动参数,配置信息等,如果想要了解可以去debezium官网查看参数的具体细节
    public DebeziumSourceFunction<T> build() {
        Properties props = new Properties();
        props.setProperty("connector.class", MySqlConnector.class.getCanonicalName());
        // hard code server name, because we don't need to distinguish it, docs:
        // Logical name that identifies and provides a namespace for the particular MySQL
        // database server/cluster being monitored. The logical name should be unique across all other
        // connectors, since it is used as a prefix for all Kafka topic names emanating from this connector.
        // Only alphanumeric characters and underscores should be used.
        props.setProperty("database.server.name", DATABASE_SERVER_NAME);
        props.setProperty("database.hostname", checkNotNull(hostname));
        props.setProperty("database.user", checkNotNull(username));
        props.setProperty("database.password", checkNotNull(password));
        props.setProperty("database.port", String.valueOf(port));
        props.setProperty("database.history.skip.unparseable.ddl", String.valueOf(true));
        // debezium use "long" mode to handle unsigned bigint by default,
        // but it'll cause lose of precise when the value is larger than 2^63,
        // so use "precise" mode to avoid it.
        props.put("bigint.unsigned.handling.mode", "precise");
        if (serverId != null) { props.setProperty("database.server.id", String.valueOf(serverId)); }
        if (databaseList != null) { props.setProperty("database.whitelist", String.join(",", databaseList)); }
        if (tableList != null) { props.setProperty("table.whitelist", String.join(",", tableList)); }
        if (serverTimeZone != null) { props.setProperty("database.serverTimezone", serverTimeZone); }

        // 判断开始消费位置,在sqlSourceBuilder中构建的参数,没有则为null
        DebeziumOffset specificOffset = null;
        switch (startupOptions.startupMode) {
            case INITIAL:
                props.setProperty("snapshot.mode", "initial");
                break;
            case EARLIEST_OFFSET:
                props.setProperty("snapshot.mode", "never");
                break;
            case LATEST_OFFSET:
                props.setProperty("snapshot.mode", "schema_only");
                break;
            case SPECIFIC_OFFSETS:
                props.setProperty("snapshot.mode", "schema_only_recovery");
                specificOffset = new DebeziumOffset();
                Map<String, String> sourcePartition = new HashMap<>();
                sourcePartition.put("server", DATABASE_SERVER_NAME);
                specificOffset.setSourcePartition(sourcePartition);
                Map<String, Object> sourceOffset = new HashMap<>();
                sourceOffset.put("file", startupOptions.specificOffsetFile);
                sourceOffset.put("pos", startupOptions.specificOffsetPos);
                specificOffset.setSourceOffset(sourceOffset);
                break;
            case TIMESTAMP:
                checkNotNull(deserializer);
                props.setProperty("snapshot.mode", "never");
                deserializer = new SeekBinlogToTimestampFilter<>(startupOptions.startupTimestampMillis, deserializer);
                break;
            default:
                throw new UnsupportedOperationException();
        }

        if (dbzProperties != null) {
            props.putAll(dbzProperties);
            // Add default configurations for compatibility when set the legacy mysql connector implementation
            if (LEGACY_IMPLEMENTATION_VALUE.equals(dbzProperties.get(LEGACY_IMPLEMENTATION_KEY))) {
                props.put("transforms", "snapshotasinsert");
                props.put("transforms.snapshotasinsert.type", "io.debezium.connector.mysql.transforms.ReadToInsertEvent");
            }
        }

        // 构建通用的cdc sourceFunction --> 基于RichSourceFunction
        return new DebeziumSourceFunction<>(deserializer, props, specificOffset, new MySqlValidator(props) // mysql校验器,版本信息,binlog是否为row等);
    }
}

上面内容主要是以构建source所需要的参数为主,具体我们进入到DebeziumSourceFunction中看看具体实现。

java 复制代码
// source代码,用于读取binlog,logminer等
// 实现RichSourceFunction完成source端代码的编写,实现CheckpointFunction用于保证容错相关的内容,实现CheckpointListener监听checkpoint的完成状态
public class DebeziumSourceFunction<T> extends RichSourceFunction<T> implements CheckpointedFunction, CheckpointListener, ResultTypeQueryable<T> {

    // ------------------------------列出一些比较重要的成员变量,不重要的忽略了------------------------------------------
    // ----------------------------------State-------------------------------------------------
    /* 主要用于状态的维护,当任务出现问题重启/手动重启后,维护的一些schema(record中的结构) 未消费的records(在queue中,后面会看到) offset等信息 */
    private transient volatile String restoredOffsetState;
    private transient ListState<byte[]> offsetState;
    private transient ListState<String> schemaRecordsState;

    // -----------------------------------Worker-----------------------------------------------
    /* 一个单线程的线程池,一个debeziumEngine(一个runnable的实现类)用于读取binlog数据 TODO 所以设计到多线程的交互 */
    private transient ExecutorService executor;
    private transient DebeziumEngine<?> engine;

    /* 一个consumer,用于从engine中读取数据的消费者,并将数据放入handover中 */
    private transient DebeziumChangeConsumer changeConsumer;

    /* 用于从handover中拿取数据 */
    private transient DebeziumChangeFetcher<T> debeziumChangeFetcher;

    /* 两个线程(source,engine)之间交互数据的一个桥梁 */
    private transient Handover handover;

    // ----------------------------------------我们主要介绍source的run方法,其他方法主要用于容错相关--------------------------------------
    @Override
    public void run(SourceContext<T> sourceContext) throws Exception {
        // TODO 用于engine执行的一些相关参数,不是终点内容,如果感兴趣可官网看看说明
        properties.setProperty("name", "engine");
        properties.setProperty("offset.storage", FlinkOffsetBackingStore.class.getCanonicalName());
        if (restoredOffsetState != null) {
            properties.setProperty(FlinkOffsetBackingStore.OFFSET_STATE_VALUE, restoredOffsetState);
        }
        properties.setProperty("include.schema.changes", "false");
        properties.setProperty("offset.flush.interval.ms", String.valueOf(Long.MAX_VALUE));
        properties.setProperty("tombstones.on.delete", "false");
        if (engineInstanceName == null) {
            engineInstanceName = UUID.randomUUID().toString();
        }
        properties.setProperty(FlinkDatabaseHistory.DATABASE_HISTORY_INSTANCE_NAME, engineInstanceName);
        properties.setProperty("database.history", determineDatabase().getCanonicalName());
        String dbzHeartbeatPrefix = properties.getProperty(Heartbeat.HEARTBEAT_TOPICS_PREFIX.name(), Heartbeat.HEARTBEAT_TOPICS_PREFIX.defaultValueAsString());
        this.debeziumChangeFetcher = new DebeziumChangeFetcher<>(sourceContext, deserializer, restoredOffsetState == null, // 是否是快照阶段或者state==null?
                dbzHeartbeatPrefix, handover);

        // 创建并配置engine相关参数
        this.engine = DebeziumEngine.create(Connect.class)
                .using(properties) // 参数
                .notifying(changeConsumer) // 配置consumer消费 engine读取的数据(binlog/历史数据)
                .using(OffsetCommitPolicy.always()) // offset的提交策略
                .using((success, message, error) -> {
                    if (success) {
                        handover.close();
                    } else {
                        handover.reportError(error);
                    }
                })
                .build();

        // 将engine任务提交到线程池中执行
        executor.execute(engine);
        debeziumStarted = true;

        // metric相关配置
        MetricGroup metricGroup = getRuntimeContext().getMetricGroup();
        // ....

        // 启动fetcher,循环去handover中拿取最新数据发送下游
        debeziumChangeFetcher.runFetchLoop();
    }
}

上面我们已经看了source.run的基本实现,他的主要处理逻辑在DebeziumChangeConsumer, DebeziumChangeFetcher, Handover中。

简单介绍三个类的作用和主要方法和参数

DebeziumChangeConsumer:用于消费engine读取的数据。

java 复制代码
/* 该类实现 DebeziumEngine.ChangeConsumer接口,实现handleBatch方法 相对比较简单, 另外两个成员方法主要是offset相关,非重点内容 */
// engine线程会调用handleBatch方法出传递引擎消费到的数据
public class DebeziumChangeConsumer implements DebeziumEngine.ChangeConsumer<ChangeEvent<SourceRecord, SourceRecord>> {
    @Override
    public void handleBatch(List<ChangeEvent<SourceRecord, SourceRecord>> events, RecordCommitter<ChangeEvent<SourceRecord, SourceRecord>> recordCommitter) {
        try {
            currentCommitter = recordCommitter;
            // 间接调用到handover的produce方法,该方法是阻塞的 嘻嘻嘻(如果有历史records未被消费则wait)
            handover.produce(events);
        } catch (Throwable e) {
            // Hold this exception in handover and trigger the fetcher to exit
            handover.reportError(e);
        }
    }
}

DebeziumChangeFetcher:循环从handover中获取consumer从engine读取的最新数据。

java 复制代码
public class DebeziumChangeFetcher<T> {
    private final SourceFunction.SourceContext<T> sourceContext;
    /* 保证数据发送和状态更新的一把锁 */
    private final Object checkpointLock;
    /* 用于将数据转化成我们自定义的类型,如json,String等 */
    private final DebeziumDeserializationSchema<T> deserialization;
    /* 下面自定义的collector */
    private final DebeziumCollector debeziumCollector;
    /* 见名知意,很好理解 */
    private final DebeziumOffset debeziumOffset;
    /* 用于存储在state offset的序列化器 */
    private final DebeziumOffsetSerializer stateSerializer;
    /* 心跳相关 */
    private final String heartbeatTopicPrefix;
    /* 是否恢复的状态,需要消费历史相关数据 */
    private boolean isInDbSnapshotPhase;
    private final Handover handover;

    public void runFetchLoop() throws Exception {
        try {
            // 读取mysql历史的数据,不要被名字所迷惑
            if (isInDbSnapshotPhase) {
                List<ChangeEvent<SourceRecord, SourceRecord>> events = handover.pollNext();

                synchronized (checkpointLock) {
                    LOG.info("Database snapshot phase can't perform checkpoint, acquired Checkpoint lock.");
                    handleBatch(events);
                    // 这里防止snapshot数据无法一次读取完毕,必须保证snapshot数据读取完毕才进入binlog的读取
                    while (isRunning && isInDbSnapshotPhase) {
                        handleBatch(handover.pollNext());
                    }
                }
                LOG.info("Received record from streaming binlog phase, released checkpoint lock.");
            }

            // 到这里表示snapshot的数据读取完毕,开始实时读取binlog数据
            while (isRunning) {
                // 具体的处理数据逻辑 pollNext会阻塞
                handleBatch(handover.pollNext());
            }
        } catch (Handover.ClosedException e) {
            // ignore
        }
    }

    private void handleBatch(List<ChangeEvent<SourceRecord, SourceRecord>> changeEvents) throws Exception {
        if (CollectionUtils.isEmpty(changeEvents)) {
            return;
        }
        this.processTime = System.currentTimeMillis();

        for (ChangeEvent<SourceRecord, SourceRecord> event : changeEvents) {
            SourceRecord record = event.value();
            // time相关基本都是metric相关内容,不必较真
            updateMessageTimestamp(record);
            fetchDelay = processTime - messageTimestamp;

            // 通过心跳机制来更新offset
            if (isHeartbeatEvent(record)) {
                synchronized (checkpointLock) {
                    debeziumOffset.setSourcePartition(record.sourcePartition());
                    debeziumOffset.setSourceOffset(record.sourceOffset());
                }
                continue;
            }

            // 根据不同的deserialization对数据做转换,可以看这个,比较容易理解StringDebeziumDeserializationSchema, 内部直接 record.toString即可,就是将debezium读取的record转换成我们想要的格式或者类型,debeziumCollector 就是下面自定义的collector,在deserialize中,会将转换完成的数据放入queue中
            deserialization.deserialize(record, debeziumCollector);

            // 判断数据是否为snapshot的最后一条数据,如果是则在这条数据之后转换到binlog的streaming流程
            if (!isSnapshotRecord(record)) {
                LOG.debug("Snapshot phase finishes.");
                isInDbSnapshotPhase = false; // runFetchLoop方法中使用
            }

            // 具体发送数据
            emitRecordsUnderCheckpointLock(debeziumCollector.records, record.sourcePartition(), record.sourceOffset());
        }
    }

    private void emitRecordsUnderCheckpointLock(Queue<T> records, Map<String, ?> sourcePartition, Map<String, ?> sourceOffset) {
        // 同步是保证数据的发送和offset的更新是安全,lock是可重入的(不懂可以百度,java基础内容)
        synchronized (checkpointLock) {
            T record;
            // 循环debeziumCollector的records队列,将队列中的数据依次发送到下游,
            while ((record = records.poll()) != null) {
                emitDelay = System.currentTimeMillis() - messageTimestamp;
                // 通过source的context对象将其发送到下游operator,这里转入了flink的处理逻辑,不再cdc代码之内
                sourceContext.collect(record);
            }
            debeziumOffset.setSourcePartition(sourcePartition);
            debeziumOffset.setSourceOffset(sourceOffset);
        }
    }

    // 心跳机制 ,用于更新offset的机制
    private boolean isHeartbeatEvent(SourceRecord record) {
        String topic = record.topic();
        return topic != null && topic.startsWith(heartbeatTopicPrefix);
    }

    // --------------------------------自定义collector-------------------------------------------------------
    private class DebeziumCollector implements Collector<T> {
        private final Queue<T> records = new ArrayDeque<>();
        @Override
        public void collect(T record) {
            // 将数据放入队列,queue会在别的地方进出列将数据发送下游
            records.add(record);
        }
    }
}

Handover:source线程和engine线程执行中数据交互桥梁。

java 复制代码
/* 这个类由两个线程访问, pollNext由debeziumFetcher调用,produce有debeziumConsumer调用,因为涉及多线程的调用,单纯的讲代码可能不容易理解,可以去复习一下java多线程知识内容,或者自己debug一下看看调用流程就比较容易理解了 */
@ThreadSafe // 表示类是线程安全的,这类涉及engine和source线程两个线程操作,内部的实现保证了线程安全
public class Handover implements Closeable {
    private static final Logger LOG = LoggerFactory.getLogger(Handover.class);
    private final Object lock = new Object();

    @GuardedBy("lock") // 注解表示该变量受lock的保护, 不是重点勿关注
    private List<ChangeEvent<SourceRecord, SourceRecord>> next;

    @GuardedBy("lock")
    private Throwable error;

    private boolean wakeupProducer;

    /* debeziumFetcher 调用,当没有数据的时候进入wait状态,wait状态的时候cpu是不会调用wait状态的线程,另一个线程就可以占用cpu的全部时间片 */
    public List<ChangeEvent<SourceRecord, SourceRecord>> pollNext() throws Exception {
        // 同步代码块才可以使用wait和notifyAll,为什么使用这种方式,因为只有两个线程,所以这种方式实现简单,如果线程多可以通过juc的lock去做或者其他方式也可以
        synchronized (lock) {
            // 没有数据没有异常则持续循环进入wait状态,为了防止虚假唤醒的情况
            while (next == null && error == null) {
                lock.wait();
            }
            List<ChangeEvent<SourceRecord, SourceRecord>> n = next;
            // 上面的循环可以退出的时候,说明一定是有数据或者有异常,不存在其他的情况
            if (n != null) {
                // 将next置为null 下面会根据此条件作为判断条件
                next = null;
                // 唤醒其他等待线程,当然只可能是engine线程
                lock.notifyAll();
                return n;
            } else {
                // 将异常抛出
                ExceptionUtils.rethrowException(error, error.getMessage());
                // 上面方法一定会抛出异常,改代码只是为了去掉编译警告...
                return Collections.emptyList();
            }
        }
    }

    public void produce(final List<ChangeEvent<SourceRecord, SourceRecord>> element) throws InterruptedException {
        checkNotNull(element);
        synchronized (lock) {
            // next不等一直进入wait状态
            while (next != null && !wakeupProducer) {
                lock.wait();
            }

            wakeupProducer = false;

            // 有异常抛出异常,没异常将接受新数据,并唤醒fetcher线程
            if (error != null) {
                ExceptionUtils.rethrow(error, error.getMessage());
            } else {
                next = element;
                lock.notifyAll();
            }
        }
    }
}

上面代码即是基于RichSourceFunction实现的cdc主要代码,其实不算难,但是前人写的代码是已经把很多问题已经考虑进入,对代码的抽象也很好,扩展起来很方便,api设计对与我们开发者来说很容易使用。

三. mysql-cdc源码 - 新Source接口的实现

1.11版本之后flink提供了新的source接口,可以提前预习一波。

FLINK-10740

简单介绍一下

  • SourceReader:对split的数据进行读取操作,比如:读取一个分区,一个块等,当然不只局限与一个分区,根据自己的实现来。

  • SplitEnumerator:负责对数据源进行切分或者发现分区等,比如:发现kafka的分区,对文件划分块等。

上述的比较简单,实际上比这复杂一点,所以在新的source接口实现一个source是比较难的事情,不过熟悉之后都一样。

提前说明:

  • 一个split我们可以认为是一个切片,mysql-cdc中,假想情况下:一张的一部分中,比如开始主键1到结束主键10,那么该split就表示这些数据,在具体读取数据的时候是有readTask来去读,那么他就会通过split标记的点位来进行数据的读取,当然一个readTask不止会执行一个split。

  • snapshot表示的是读取数据库的历史全量数据。

  • binlog表示当我们snapshot阶段结束后开始binlog阶段,即我们开始读取的binlog数据了。

  • 先执行snapshot阶段,后执行binlog阶段。

代码的生成和旧版是相同的,只不过是内部执行的逻辑存在变化,新的source接口实现的cdc代码比较复杂,涉及的内容比较多,可能比较晕,后面自己可以根据源码debug走一走。

由于代码过多,主要讲解重点的内容,不重要的跳过了。

java 复制代码
// 实现了两个接口 source,和 resultTypeQueryable(比较简单就一个获取结果类型信息的接口) , 主要代码还是在source接口的实现
// T 为输出类型,MySqlSplit是mysql的分割器,PendingSplitsState表示Enumerator的状态对象
public class MySqlSource<T> implements Source<T, MySqlSplit, PendingSplitsState>, ResultTypeQueryable<T> {

    private final MySqlSourceConfigFactory configFactory;
    private final DebeziumDeserializationSchema<T> deserializationSchema;

    /* 通过构造者模式构建source所需要的参数,简单说明一下,里面的参数,通过MySqlSourceConfigFactory添加参数,在build方法中,将factory作为参数构建出MySqlSource
    -------------------------------------讲解一下对应关系------------------------------------------------
    MySqlSourceConfigFactory 可以根据不同的subtask创建对应的MySqlSourceConfig
    MySqlSourceConfig 可以构建 MySqlConnectorConfig
    MySqlConnection 通过 DebeziumUtil.createMySqlConnection(mySqlSourceConfig.getDbzConfiguration())方法构建
    上面的一个config比较混乱,名字也比较不容易理解,后面用到的时候会简单提一下,这里主要是有一个印象,不要被一些配置搞混
    */
    public static <T> MySqlSourceBuilder<T> builder() {
        return new MySqlSourceBuilder<>();
    }

    // 由MySqlSourceBuilder.build方法创建
    MySqlSource(MySqlSourceConfigFactory configFactory, DebeziumDeserializationSchema<T> deserializationSchema // 与老版source的deserialization一样) {
        this.configFactory = configFactory;
        this.deserializationSchema = deserializationSchema;
    }

    @Override // 流批一体的source,表示有界性,新source接口的特性
    public Boundedness getBoundedness() { return Boundedness.CONTINUOUS_UNBOUNDED; }

    /* 构建sourceReader */
    @Override
    public SourceReader<T, MySqlSplit> createReader(SourceReaderContext readerContext) throws Exception {
        // 前面提到了,根据subtask索引创建对应的config
        MySqlSourceConfig sourceConfig = configFactory.createConfig(readerContext.getIndexOfSubtask());
        // 一个阻塞队列,多线程交互用的,不必深入
        FutureCompletingBlockingQueue<RecordsWithSplitIds<SourceRecord>> elementsQueue = new FutureCompletingBlockingQueue<>();
        // metric相关
        final MySqlSourceReaderMetrics sourceReaderMetrics = new MySqlSourceReaderMetrics(readerContext.metricGroup());
        sourceReaderMetrics.registerMetrics();
        // 通过supplier函数构建一个SplitReader,解耦的作用,主要看里面的MySqlSplitReader实现即可
        Supplier<MySqlSplitReader> splitReaderSupplier = () -> new MySqlSplitReader(sourceConfig, readerContext.getIndexOfSubtask());

        // 构建了一个具体的sourceReader
        return new MySqlSourceReader<>(elementsQueue, splitReaderSupplier, new MySqlRecordEmitter<>(deserializationSchema, sourceReaderMetrics, sourceConfig.isIncludeSchemaChanges()), readerContext.getConfiguration(), readerContext, sourceConfig);
    }

    @Override
    public SplitEnumerator<MySqlSplit, PendingSplitsState> createEnumerator(SplitEnumeratorContext<MySqlSplit> enumContext) {
        // 因为只会生成一次所以生成一个sourceConfig即可
        MySqlSourceConfig sourceConfig = configFactory.createConfig(0);
        // 检验mysql
        final MySqlValidator validator = new MySqlValidator(sourceConfig);
        validator.validate();

        final MySqlSplitAssigner splitAssigner;
        // 判断开始条件如果是initial则先读取mysql table的数据(代码中叫做snapshot),然后再继续读取binlog的数据,如果不是initial状态,则直接从binlog开始读取
        if (sourceConfig.getStartupOptions().startupMode == StartupMode.INITIAL) {
            try (JdbcConnection jdbc = openJdbcConnection(sourceConfig)) {
                final List<TableId> remainingTables = discoverCapturedTables(jdbc, sourceConfig);
                boolean isTableIdCaseSensitive = DebeziumUtils.isTableIdCaseSensitive(jdbc);
                splitAssigner = new MySqlHybridSplitAssigner(sourceConfig, enumContext.currentParallelism(), remainingTables, isTableIdCaseSensitive);
            } catch (Exception e) {
                throw new FlinkRuntimeException("Failed to discover captured tables for enumerator", e);
            }
        } else {
            // 之有binlog的split逻辑
            splitAssigner = new MySqlBinlogSplitAssigner(sourceConfig);
        }
        // 创建对应发的SplitEnumerator,用于构建split给reader读取
        return new MySqlSourceEnumerator(enumContext, sourceConfig, splitAssigner);
    }

    // 恢复SplitEnumerator,比如任务故障重启,会根据不同的checkpoint恢复SplitEnumerator,用于继续之前未完成的读取操作
    @Override
    public SplitEnumerator<MySqlSplit, PendingSplitsState> restoreEnumerator(SplitEnumeratorContext<MySqlSplit> enumContext, PendingSplitsState checkpoint) {
        MySqlSourceConfig sourceConfig = configFactory.createConfig(0);

        final MySqlSplitAssigner splitAssigner;
        if (checkpoint instanceof HybridPendingSplitsState) {
            splitAssigner = new MySqlHybridSplitAssigner(sourceConfig, enumContext.currentParallelism(), (HybridPendingSplitsState) checkpoint);
        } else if (checkpoint instanceof BinlogPendingSplitsState) {
            splitAssigner = new MySqlBinlogSplitAssigner(sourceConfig, (BinlogPendingSplitsState) checkpoint);
        } else {
            throw new UnsupportedOperationException("Unsupported restored PendingSplitsState: " + checkpoint);
        }
        return new MySqlSourceEnumerator(enumContext, sourceConfig, splitAssigner);
    }

    // ------------------容错相关,不是重点-----------------
    @Override
    public SimpleVersionedSerializer<MySqlSplit> getSplitSerializer() { return MySqlSplitSerializer.INSTANCE; }

    @Override
    public SimpleVersionedSerializer<PendingSplitsState> getEnumeratorCheckpointSerializer() { return new PendingSplitsStateSerializer(getSplitSerializer()); }

    // 返回值类型的提取
    @Override
    public TypeInformation<T> getProducedType() { return deserializationSchema.getProducedType(); }
}

上面的代码中我们可以看到source的实现,主要是构建sourceReader和splitEnumerator,以及容错内容,相关的处理逻辑也封装在相应的对象中,下面我们对其内部逐步剖析。

java 复制代码
/* 在看其他内容之前,我们可以看看如何对mysql进行split操作,在snapshot是通过主键来split的,binlog的只从当前offset位置开始消费,
这里是混合的一个split,另外还存在binlog和snapshot的splitAssigner,不过我们根据主要看看大致逻辑,具体到某一直可以自己阅读理解,
解释一下 : 先读取mysql历史数据即snapshot阶段,然后再进行当前mysql-binlog的位置开始消费,所以这个混合的意义就是先读取全量数据,然后从最新的binlog开始读取,完成cdc读取数据的过程 */
public class MySqlHybridSplitAssigner implements MySqlSplitAssigner {
    private final int splitMetaGroupSize;
    private boolean isBinlogSplitAssigned;
    private final MySqlSnapshotSplitAssigner snapshotSplitAssigner;

    public MySqlHybridSplitAssigner(MySqlSourceConfig sourceConfig, int currentParallelism, List<TableId> remainingTables, boolean isTableIdCaseSensitive) {
        this(new MySqlSnapshotSplitAssigner(sourceConfig, currentParallelism, remainingTables, isTableIdCaseSensitive), false, sourceConfig.getSplitMetaGroupSize());
    }

    public MySqlHybridSplitAssigner(MySqlSourceConfig sourceConfig, int currentParallelism, HybridPendingSplitsState checkpoint) {
        this(new MySqlSnapshotSplitAssigner(sourceConfig, currentParallelism, checkpoint.getSnapshotPendingSplits()), checkpoint.isBinlogSplitAssigned(), sourceConfig.getSplitMetaGroupSize());
    }

    private MySqlHybridSplitAssigner(MySqlSnapshotSplitAssigner snapshotSplitAssigner, boolean isBinlogSplitAssigned, int splitMetaGroupSize) {
        this.snapshotSplitAssigner = snapshotSplitAssigner;
        this.isBinlogSplitAssigned = isBinlogSplitAssigned;
        this.splitMetaGroupSize = splitMetaGroupSize;
    }

    @Override
    public void open() {
        snapshotSplitAssigner.open();
    }

    // 主要返回下一个split,没有则返回一个空, optional可以jdk8的新特性,用于解决空指针的一个类
    @Override
    public Optional<MySqlSplit> getNext() {
        // 下面的方法可以见名知意,自行理解即可
        if (snapshotSplitAssigner.noMoreSplits()) {
            if (isBinlogSplitAssigned) {
                return Optional.empty();
            } else if (snapshotSplitAssigner.isFinished()) { // 当snapshot完成后,开始binlog的split流程
                // we need to wait snapshot-assigner to be finished before
                // assigning the binlog split. Otherwise, records emitted from binlog split
                // might be out-of-order in terms of same primary key with snapshot splits.
                isBinlogSplitAssigned = true;
                return Optional.of(createBinlogSplit());
            } else {
                // binlog split is not ready by now
                return Optional.empty();
            }
        } else {
            // snapshot assigner still have remaining splits, assign split from it
            return snapshotSplitAssigner.getNext();
        }
    }

    // splitAssigner是否在等待已完成split回调,即onFinishedSplits
    @Override
    public boolean waitingForFinishedSplits() {
        return snapshotSplitAssigner.waitingForFinishedSplits();
    }

    // 获取已完成的split并且包含他的元数据,可以根据已经完成snapshot(snapshot的某一个split)生成对应binlog的split
    @Override
    public List<FinishedSnapshotSplitInfo> getFinishedSplitInfos() {
        return snapshotSplitAssigner.getFinishedSplitInfos();
    }

    // 使用已完成的binlog偏移量来处理已完成的split,用于确定何时生成binlog split以及生成什么binlog split,就是回调
    @Override
    public void onFinishedSplits(Map<String, BinlogOffset> splitFinishedOffsets) {
        snapshotSplitAssigner.onFinishedSplits(splitFinishedOffsets);
    }

    // 向此splitAssigner添加一组split,当某些split处理失败,则需要重新添加分割时调用此方法
    @Override
    public void addSplits(Collection<MySqlSplit> splits) {
        List<MySqlSplit> snapshotSplits = new ArrayList<>();
        for (MySqlSplit split : splits) {
            if (split.isSnapshotSplit()) {
                snapshotSplits.add(split);
            } else {
                // we don't store the split, but will re-create binlog split later
                isBinlogSplitAssigned = false;
            }
        }
        snapshotSplitAssigner.addSplits(snapshotSplits);
    }

    // ----------------------------checkpoint 容错相关----------------------------------------
    @Override
    public PendingSplitsState snapshotState(long checkpointId) {
        return new HybridPendingSplitsState(snapshotSplitAssigner.snapshotState(checkpointId), isBinlogSplitAssigned);
    }

    @Override
    public void notifyCheckpointComplete(long checkpointId) {
        snapshotSplitAssigner.notifyCheckpointComplete(checkpointId);
    }

    @Override
    public void close() {
        snapshotSplitAssigner.close();
    }

    // -------------------------------------binlog split部分-------------------------------------------
    // 构建binlog split, 就是根据已经完成snapshot split来构建binlog split的一个过程,split代码比较简单可以自行阅读
    // 简单介绍一下 就是描述binlog的split,snapshot的split相关内容,比如snapshot,会按照主键去做split,已经table的schemas相关信息
    private MySqlBinlogSplit createBinlogSplit() {
        final List<MySqlSnapshotSplit> assignedSnapshotSplit = snapshotSplitAssigner.getAssignedSplits().values().stream().sorted(Comparator.comparing(MySqlSplit::splitId)).collect(Collectors.toList());
        Map<String, BinlogOffset> splitFinishedOffsets = snapshotSplitAssigner.getSplitFinishedOffsets();
        final List<FinishedSnapshotSplitInfo> finishedSnapshotSplitInfos = new ArrayList<>();
        BinlogOffset minBinlogOffset = null;
        for (MySqlSnapshotSplit split : assignedSnapshotSplit) {
            // find the min binlog offset
            BinlogOffset binlogOffset = splitFinishedOffsets.get(split.splitId());
            if (minBinlogOffset == null || binlogOffset.isBefore(minBinlogOffset)) {
                minBinlogOffset = binlogOffset;
            }
            finished
end