ITPUB

实战复盘:Flink CDC 整库同步 MySQL 到 Doris的 DDL(create) 咋实现?

看过我往期文章的同学应该知道,对于用 Flink CDC 来同步上游数据这个事,我是真的没少折腾,其中研究出的一些解决方案,已经为一些公司或者项目,实实在在解决了很多生产问题。

而这一次,难度又升级了,有客户要求,需要在用 Flink CDC Source 同步 MySQL 整库到 Doris 的情况下,如果上游有新建表,也必须把新增的表结构,也一并给「实时同步」过来。

可别小看这「实时」两字,从实现上来说,一下子就上了难度。

还记得之前,给一个项目搞定他们 PG 整库同步 Doris 的方案中,如果上游新增表,收到通知后,下游可以根据我提供的一个,表结构同步命令,执行一下,然后,再重启 CDC 任务(详见之前历史文章),算是对这个场景的交差。

但这一次,不一样了,对人家上游业务库来说,「随时可能新建表」,而且,「我不会通知你」,但你依然还得满足「实时同步」,果然,甲方就是硬气。

看在钱的份上,虽然心里骂骂咧咧,但嘴上必须说着「没问题,包搞定」。

0. 可行方案

先说这里的最难点,是上游 MySQL 的表结构,跟下游 Doris 表结构的映射关系,怎么在 Flink 的 coding 中体现。

如果我一个个去根据经验,敲字段间的映射关系,效率低不说,估计累个半死效果还不一定好,明显不划算。

对于这一点,其实可以借助 Doris 官网中描述的,用 Flink pipeline 整库同步 MySQL 时的案例。

来,看这里:

Image

这里提到了个核心关键类「CdcTools」,对,就是这个货。

既然它可以直接同步整库的表结构,以及每张表里的数据,那么我们利用它的功能,是不是就可以达到「同步建表」的目的?

答案必须是的。

那接下来咱们的重点,就是看怎么把这个类的相关能力,在 Flink CDC Source 的编码环境里体现出来。

1. CdcTools 调用方式

既然这个类在命令行中能用,那在 Flink 的代码框架里,肯定一样可以用,怎么玩?

来,我已经替你们试过了,可以直接用这种方式,调用这个类:

CdcTools.main(Array[String](
"mysql-sync-database",
"--create-table-only",
"--database", "test",
"--mysql-conf", "hostname=192.xxx.xx.xxx",
"--mysql-conf", "port=3306",
"--mysql-conf", "username=mysql_user",
"--mysql-conf", "password=****",
"--mysql-conf", "database-name=test",
"--sink-conf", "fenodes=192.168.xxx.xx:8030",
"--sink-conf", "username=doris_user",
"--sink-conf", "password=****",
"--sink-conf", "jdbc-url=jdbc:mysql://192.168.xxx.xxx:9030",
"--sink-conf", "sink.label-prefix=label01",
"--table-conf", "replication_num=2"))

这样,就可以把上游 msyql DDL 中的「新建表」,给同步到下游的 Doris 库里。

是不是很简单?但前提是,你得把这个类的相关源码,都给务必读懂才好使。

2. Flink CDC 实现策略

知道了「新建表」的同步方式之后,接下来,我们就需要明确一件事——「新建表」跟「数据同步」,是两个不同性质的动作。

怎么理解?

普通的「数据同步」,其实是 insert、update、delete,对于这个,Doris 提供了专门的写 API。

但「新建表」则不一样,它需要我们换一种策略,去调用刚才的 CdcTools。

所以,这两种情况,就需要「分开处理」,当然,如果你能一次性处理,并且还能保证效率的话,那说明,你比我厉害。

至于怎么分开处理?最常见的思路就是——用侧输出流。

对于 create (新建表)操作,用侧输出流;

而对于普通数据的同步,直接用主输出流。

3. Flink 核心编码实现

先看一眼 msyql 对于 DDL 中的 create 操作,收到的原始数据,长这样:

{"source":{"server_id":1,"version":"1.9.8.Final","file":"binlog.000026","connector":"mysql","pos":1027413231,"name":"mysql_binlog_source","row":0,"ts_ms":1763089434518,"snapshot":"false","db":"test01","table":"t06"},"historyRecord":"{\"source\":{\"file\":\"binlog.000026\",\"pos\":1027413231,\"server_id\":1},\"position\":{\"transaction_id\":null,\"ts_sec\":1763089434,\"file\":\"binlog.000026\",\"pos\":1027413369,\"server_id\":1},\"databaseName\":\"test01\",\"ddl\":\"alter table t06 add column addr text after age\",\"tableChanges\":[{\"type\":\"ALTER\",\"id\":\"\\\"test01\\\".\\\"t06\\\"\",\"table\":{\"defaultCharsetName\":\"utf8mb4\",\"primaryKeyColumnNames\":[\"ID\"],\"columns\":[{\"name\":\"ID\",\"jdbcType\":4,\"typeName\":\"INT\",\"typeExpression\":\"INT\",\"charsetName\":null,\"position\":1,\"optional\":false,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":false,\"enumValues\":[]},{\"name\":\"NAME\",\"jdbcType\":12,\"typeName\":\"TEXT\",\"typeExpression\":\"TEXT\",\"charsetName\":\"utf8mb4\",\"position\":2,\"optional\":false,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":false,\"enumValues\":[]},{\"name\":\"age\",\"jdbcType\":4,\"typeName\":\"INT\",\"typeExpression\":\"INT\",\"charsetName\":null,\"position\":3,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"defaultValueExpression\":\"99\",\"enumValues\":[]},{\"name\":\"addr\",\"jdbcType\":12,\"typeName\":\"TEXT\",\"typeExpression\":\"TEXT\",\"charsetName\":\"utf8mb4\",\"position\":4,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"enumValues\":[]}]},\"comment\":null}]}"}

这个 josn 里面有个关键 key——historyRecord。

所以,对于这部分的数据判断,就如下面这样:

if (json.containsKey("historyRecord")){//DDL 侧流
val histJson = json.getJSONObject("historyRecord")
if (histJson.containsKey("tableChanges") && !histJson.getJSONArray("tableChanges").isEmpty){
val tableChange = histJson.getJSONArray("tableChanges").getJSONObject(0)
val ddlType = tableChange.getString("type") // ALTER/CREATE/DROP
if ("CREATE".equalsIgnoreCase(ddlType)){
              ctx.output(ddlSideOutput, histJson)
          }
      }
  }

然后,把它分到「create 侧流」里。

至于非 create 的 insert、update、delete,就放到「insert 主流」里。

接下来,为了提高这个 create 操作的效率,这里考虑对这个侧流,利用「异步 IO」来处理:

AsyncDataStream.unorderedWait(createDS, newAsyncFunction[JSONObject, String] {//异步执行DDL操作
overridedefasyncInvoke(input: JSONObject, resultFuture: async.ResultFuture[String]): Unit = {
val tableChange = input.getJSONArray("tableChanges").getJSONObject(0)
val ddlType = tableChange.getString("type") // ALTER/CREATE/DROP
if ("CREATE".equalsIgnoreCase(ddlType)){             
//执行 create 操作
CdcTools.main(Array[String](
"mysql-sync-database",
"--create-table-only",
"--database", "test",
"--mysql-conf", "hostname=192.xx.xxx.xxx",
"--mysql-conf", "port=3306",
"--mysql-conf", "username=mysql_user",
"--mysql-conf", "password=***",
"--mysql-conf", "database-name=test",
"--sink-conf", "fenodes=192.168.xxx.xxx:8030",
"--sink-conf", "username=doris_user",
"--sink-conf", "password=****",
"--sink-conf", "jdbc-url=jdbc:mysql://192.168.xxx.xxx:9030",
"--sink-conf", "sink.label-prefix=label01",
"--table-conf", "replication_num=2"))
                }
            }
        }, 1000, TimeUnit.MILLISECONDS, 1).setParallelism(1)

而对于「主流」的处理方式,就非常简单了——该怎么着就怎么着。

因为没有什么新鲜东西,所以这里对于「主流」的处理逻辑就不赘述,有兴趣可以参考之前的文章:

Doris 整库同步生产表,终于搞定了!

4. 有个值得注意的地方

其实这个所谓的「注意」,现在已经是被熟知的潜规则了,那就是对于「主流」来说,除了必要的业务字段外,还必须得添加一个额外的「Doris 删除标记字段」。

就是这个:

rowJson.put("__DORIS_DELETE_SIGN__", 0)

虽然对于很多人来说有点莫名其妙,但你要是不加,人家是会报错滴:

Reason: There is no column matching jsonpaths in the json file, columns:[id, name, __DORIS_DELETE_SIGN__], please check columns and jsonpaths:. src line {"name":"X","id":5}

最后

其实对于同步 MySQL 上游 DDL 中的 create 操作,本身不能说有多大的难度,但如果你想自己来实现 MySQL 表,到 Doris 表的映射过程,那确实老麻烦了(一开始尝试过,后面放弃了)。

但如果能巧妙的利用 flink-doris-connector 包里,这个 CdcTools 工具类自带的一些现成解决方案,不失为一种聪明的办法。

当然,DDL 除了 create 外,还有 alter 操作,那么下一期,咱一起再来看看,针对 alter 的情况,又该怎么来玩。

Image