ITPUB

实战复盘:Flink CDC 整库同步 MySQL 到 Doris,ddl 中的 alter 咋整?

如果你用过 Flink CDC 的 pipeline,会发现,它是目前为止,用 CDC 的方式,实时数据同步「最简单」的解决方案,甚至没有之一。

理论上(或者运气足够好),仅用这种「命令行 + yaml 配置文件」的玩法,就可以「整库」同步支持的「上下游」数据源,比如,从 MySQL 整库到 Doris。

但,在落地实践的时候,真实的业务需求,往往对我们提出了更高的要求,比如客户指着你的鼻子说:

——老子除了要实时同步 CDC 数据外,还必须要实时同步业务库的 DDL 操作。

你谄媚一笑:

——包的。

根据客户的描述,这里的业务库 DDL,指的是两类操作—— create + alter。

而这些要求,根据以往的实测经验,很明显,用 Flink CDC pipeline 是打死都搞不定的,只能上 CDC Source,手动去撸。

而对于 DDL 中的 create 操作,上文已经提供了一种,经过我实测之后可行的解决方案。

实战复盘:Flink CDC 同步 MySQL 到 Doris,DDL 到底该怎么配才不出错?

那么这一篇,咱们继续来看 DDL 中如果遇到 alter,又该怎么玩?

0. 实现方式

其实相比同步上次的 create 操作,根据个人经验,这个同步 alter,反而要更简单一些,因为它,相对来说更容易想到。

上次的同步 create 操作,因为涉及到 2 张表之间的「表结构映射关系」,而 MySQL 跟 Doris 的建表语法,又差异比较大,所以它们之间的对应关系,如果纯自己手动撸,我肯定不太愿意。

刚好 flink-doris-connector 有个类 CdcTools,专门提供了这种上下游的「新建表同步」功能,所以在 Flink 代码中,巧妙的利用它,就能搞定这个上游 create 同步问题。

这次咱们的 alter 操作,目前我暂时没有找到类似的工具类玩法,只能靠自己写,但好在,这部分的操作,足够简单。

为什么这么说?

因为这个客户要求的 alter 操作,其实只包含一种情况——添加新字段。

而且还有个好消息,这个 MySQL 的 alter 语法,跟 Doris 的 alter 语法是兼容的。

也就是说,大部分情况下(除了少数字段类型不一致外),你可以直接拿 MySQL 的 alter SQL,到 Doris 上执行,没有问题。

好,就冲这一点,咱们的解决思路,就有了。

1. 具体怎么做

首先,我们需要把这个 alter 操作,给「识别」出来,Flink CDC 把它读到之后,是这样的一条原始 json:

{"source":{"version":"1.9.8.Final","connector":"mysql","name":"mysql_binlog_source","ts_ms":1763538523213,"snapshot":"false","db":"test","sequence":null,"table":"test_dist","server_id":1,"gtid":null,"file":"binlog.000026","pos":1033367470,"row":0,"thread":null,"query":null},"historyRecord":"{\"source\":{\"file\":\"binlog.000026\",\"pos\":1033367470,\"server_id\":1},\"position\":{\"transaction_id\":null,\"ts_sec\":1763538523,\"file\":\"binlog.000026\",\"pos\":1033367609,\"server_id\":1},\"databaseName\":\"test\",\"ddl\":\"alter table test_dist add column c03 int after addr\",\"tableChanges\":[{\"type\":\"ALTER\",\"id\":\"\\\"test\\\".\\\"test_dist\\\"\",\"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\":\"VARCHAR\",\"typeExpression\":\"VARCHAR\",\"charsetName\":\"utf8mb4\",\"length\":50,\"position\":2,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"enumValues\":[]},{\"name\":\"age\",\"jdbcType\":4,\"typeName\":\"INT\",\"typeExpression\":\"INT\",\"charsetName\":null,\"position\":3,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"enumValues\":[]},{\"name\":\"addr\",\"jdbcType\":12,\"typeName\":\"TEXT\",\"typeExpression\":\"TEXT\",\"charsetName\":\"utf8mb4\",\"position\":4,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"enumValues\":[]},{\"name\":\"c03\",\"jdbcType\":4,\"typeName\":\"INT\",\"typeExpression\":\"INT\",\"charsetName\":null,\"position\":5,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"enumValues\":[]},{\"name\":\"c02\",\"jdbcType\":4,\"typeName\":\"INT\",\"typeExpression\":\"INT\",\"charsetName\":null,\"position\":6,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"enumValues\":[]},{\"name\":\"c01\",\"jdbcType\":4,\"typeName\":\"INT\",\"typeExpression\":\"INT\",\"charsetName\":null,\"position\":7,\"optional\":true,\"autoIncremented\":false,\"generated\":false,\"comment\":null,\"hasDefaultValue\":true,\"enumValues\":[]}]},\"comment\":null}]}"}

观察它的特点,跟上次那个 create 不一样的地方在于,它的 ddlType,是 alter。

接着,我们需要用 Flink CDC 同时处理「业务变更数据」,以及「DDL数据」,但因为它两的「性质」不一样,所以得走两条不同的线。

怎么走?上篇文章其实说了,用侧流输出的方式。

DDL 流,输出到「侧流」;

业务变更数据流,输出到「主流」。

1.1 数据分流

所以 Flink 拿到数据的第一步,就是识别出 alter 流,跟业务数据流,分流逻辑是这样的:

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 ("ALTER".equalsIgnoreCase(ddlType)){
            ctx.output(ddlSideOutput, histJson)
        }
    }
}
else{//主流
    out.collect(json)
}

1.2 对 alter 流处理

拿到 alter 侧流之后,就要提取到其中的 alter SQL,然后去 Doris 执行。

考虑到它跟「业务变更数据」不一样,不能通过 Doris 官方推荐的写 API,目前我能想到最好的方式—— JDBC 写入。

因为这里涉及到跟 Doris 的连接问题,所以,可以用 Flink 的富函数来实现(为什么没用上次的异步 IO,文末会解释):

.addSink(newRichSinkFunction[JSONObject]() {
privateval url = "jdbc:mysql://192.168.xxx.xxx:9030/test"
privateval user = "user"
privateval password = "****"
privatevar conn: Connection = null

overridedefopen(parameters: Configuration): Unit = {
Class.forName("com.mysql.jdbc.Driver")
                    conn = DriverManager.getConnection(url, user, password)
                }

overridedefinvoke(value: JSONObject, context: SinkFunction.Context): Unit = {
val tableChange = value.getJSONArray("tableChanges").getJSONObject(0)
val ddlType = tableChange.getString("type")
val ddlSql = value.getString("ddl")

val ps = conn.prepareStatement(ddlSql)
                        ps.executeUpdate()
                    }                   
            }).setParallelism(1)

当然,这里暂时只考虑 MySQL 字段类型,跟 Doris 类型完全兼容的情况(大部分常用的都兼容)。

1.3 关于主流

就用最普通的,Doris 推荐的写方式就好了。

同样,因为这部分确实没什么新鲜东西,就不展开了(有兴趣可以参考之前的历史文章)。

2. 值得注意的地方

上篇文章提到了,用 Flink CDC Source 这种方式整库同步到 Doris 时,必须要在业务字段里,额外加上一个「Doris 删除标记」字段。

其实在这个前,还漏掉了一个很重要的信息点,那就是对于上游(不止 MySQL)的 DDL 操作,想要被识别到,必须得打开几个开关。

否则,上游 DDL 的动静闹得再大,Flink 程序默认是不知道的。

你得像下面这么来设置:

Image

最后

认真看过上一篇文章的同学,可能有这样的怀疑:咦? 上次的 DDL create 操作,不是用的「异步 IO 函数」搞定的吗?怎么这次不是了呢?

原因很简单,当我用它来测试这次的 alter 操作时发现,Flink 程序会有莫名其妙的报错,所以,就换了。

经过这次的整库同步需求,我发现,数据同步(迁移)这件事,想要玩明白、玩深入,还真不是一件容易的事情。

你觉得呢?

Image