万字长文Flinkcdc源码精讲推荐收藏
编者荐语
相对比较全面的一篇文章。
以下文章来源于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接口,可以提前预习一波。
简单介绍一下
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
