Apache Paimon 数据出仓源码导读(四):从零跑通 Paimon 到 MySQL:环境、建表与第一条数据
2026/9/4 6:03:43 网站建设 项目流程

上一篇,我们从源码角度总览了 Paimon 出仓到 MySQL 时的 UPSERT、DELETE、批量写入和故障重放。

不过,如果还没有亲手跑过这条链路,一上来就看RowKindJdbcOutputFormatexecuteBatch()和 Checkpoint,很容易把几个层次混在一起。

所以从这一篇开始,我们暂时放慢速度,把“Paimon 出仓到 MySQL”拆成几篇独立文章。

这一篇只完成一件事:

从零搭好一条最小链路,亲眼看到 Paimon 中的一条订单在 MySQL 中新增、更新和删除。

先不追每一个源码调用,也先不讨论失败重放。等链路真正跑通以后,下一篇再把“一条记录怎样进入 JDBC Sink、怎样在 Buffer 中攒成 Batch”逐层拆开。

一、这次到底要搭一条什么链路

最终运行的是一条持续存在的 Flink SQL 作业:

Paimon 主键表 orders ↓ Paimon Source 持续读取当前全量和后续变化 ↓ Flink SQL:INSERT INTO mysql_orders SELECT ... ↓ JDBC Sink ↓ MySQL 物理表 orders_rt

这里有三张“表”,但它们并不都是同一种东西:

名字存在哪里它是什么
ordersPaimon真正保存订单数据的湖表
mysql_ordersFlink Catalog指向 MySQL 的 JDBC 表定义,本身不保存另一份数据
orders_rtMySQL真正接收出仓结果的物理表

很多初学者会误以为执行CREATE TABLE mysql_orders后,Flink 会自动去 MySQL 建表。实际上它只是向 Flink 注册了一份连接和字段映射。

MySQL 中的orders_rt仍然需要提前创建。

二、运行前需要准备哪些组件

本文环境基于:

Apache Paimon 1.4.2 Apache Flink 1.20.1 Flink JDBC Connector 3.3.0-1.20 MySQL 8.x

除了 Flink 和 Paimon,还需要两个容易漏掉的依赖:

依赖解决什么问题缺失时常见现象
Flink JDBC Connector让 Flink 认识'connector' = 'jdbc'找不到jdbcFactory
MySQL Connector/J让 JDBC 能真正连接 MySQL找不到 MySQL Driver 或无法创建连接

Flink 的二进制发行包通常不自带 JDBC Connector 和数据库 Driver。对于普通 Standalone 集群,最直观的做法是让 JobManager 和所有 TaskManager 的 Flinklib目录都包含版本匹配的 JAR,然后重启对应进程。

为什么不能只放在 SQL Client 所在机器?

因为 SQL Client 主要负责解析和提交作业,真正执行 JDBC Sink 的是 TaskManager。TaskManager 看不到 Driver,作业仍然会失败。

先做三个检查

第一,确认 Paimon 和 Flink 大版本匹配。不要把面向 Flink 1.18 的 Bundle 直接放进 Flink 1.20。

第二,确认 JDBC Connector 的后缀与 Flink 版本匹配。本文使用的是3.3.0-1.20

第三,确认 MySQL 网络可达。TaskManager 所在机器必须能访问 MySQL Host 和端口,不是只有你的电脑能连上就够了。

三、先在 MySQL 建真正的目标表

先创建实验数据库:

CREATEDATABASEIFNOTEXISTSservingDEFAULTCHARACTERSETutf8mb4;

再创建目标表:

CREATETABLEserving.orders_rt(idBIGINTNOTNULL,amountDECIMAL(10,2),statusVARCHAR(32),PRIMARYKEY(id))ENGINE=InnoDB;

这条PRIMARY KEY (id)不是可有可无的装饰。后面同一个订单再次到来时,MySQL 正是依靠这个真实主键判断:

id 不存在 -> 插入新行 id 已存在 -> 更新原来的行

创建完成后,不要只凭印象判断主键已经存在,直接检查:

SHOWCREATETABLEserving.orders_rt;

结果中应该能看到类似:

PRIMARYKEY(`id`)

写入账号需要哪些权限

JDBC Sink 至少会执行 INSERT、UPDATE 语义和 DELETE,因此写入账号要具备目标表对应权限。

生产环境中建议由 DBA 创建最小权限账号,不要为了省事给全库管理员权限。本文 DDL 中继续使用占位账号:

username = flink_writer password = ******

不要把真实密码直接提交到 Git、文章截图或公开的 Flink SQL 文件中。

四、在 Flink SQL Client 中准备实验环境

进入 Flink SQL Client 后,先使用流模式:

SET'execution.runtime-mode'='streaming';

为了让第一次实验更容易观察,可以暂时把默认并行度设为 1,并开启 10 秒 Checkpoint:

SET'parallelism.default'='1';SET'execution.checkpointing.interval'='10s';

并行度设为 1 只是为了让初学实验中的日志、连接和数据顺序更容易看懂,不是生产推荐值。

生产环境的并行度要结合 MySQL 连接数、写入吞吐、热点 Key、锁等待和 Checkpoint 时长重新评估。

五、创建 Paimon 订单主键表

在 Flink SQL 中创建 Paimon 表:

CREATETABLEorders(idBIGINT,amountDECIMAL(10,2),statusSTRING,PRIMARYKEY(id)NOTENFORCED)WITH('connector'='paimon','path'='hdfs:///warehouse/orders','changelog-producer'='input');

逐项解释:

定义含义
id订单唯一标识
amount订单金额,保留两位小数
status订单状态
PRIMARY KEY (id)相同id表示同一份订单状态
NOT ENFORCEDFlink 不负责运行时检查唯一性
pathPaimon 表在文件系统中的存储位置
changelog-producer = input保存上游输入变化,供下游增量消费

这里先使用input,是为了后面能更直观地观察变化怎样传到 JDBC Sink。不同 Changelog Producer 的差别已经在上一篇单独讲过。

先写两条初始数据

执行:

INSERTINTOordersVALUES(1001,CAST(80.00ASDECIMAL(10,2)),'CREATED'),(1002,CAST(50.00ASDECIMAL(10,2)),'CREATED');

作业完成后,可以在用于查询快照的 SQL Client 会话中切到 Batch 模式,再查询 Paimon 当前结果:

SET'execution.runtime-mode'='batch';SELECTid,amount,statusFROMordersORDERBYid;

预期当前结果是:

idamountstatus
100180.00CREATED
100250.00CREATED

这里看到的是 Paimon 表的当前状态,不是历史变化流水。

六、在 Flink 中注册 MySQL JDBC Sink

现在创建 Flink JDBC 表:

CREATETABLEmysql_orders(idBIGINT,amountDECIMAL(10,2),statusSTRING,PRIMARYKEY(id)NOTENFORCED)WITH('connector'='jdbc','url'='jdbc:mysql://mysql-host:3306/serving','table-name'='orders_rt','username'='flink_writer','password'='******','sink.buffer-flush.max-rows'='100','sink.buffer-flush.interval'='1s','sink.max-retries'='3');

这一段最容易产生两个误解。

误解一:mysql_orders是又创建了一张 MySQL 表

不是。

mysql_orders是 Flink 中的逻辑表名,table-name = orders_rt才是 MySQL 里的真实物理表。

以后在 Flink SQL 中写:

INSERTINTOmysql_orders...

Connector 才会根据 URL 和table-name找到:

serving.orders_rt

误解二:Flink DDL 的主键会替 MySQL 创建约束

也不会。

Flink DDL 中的:

PRIMARYKEY(id)NOTENFORCED

告诉 Planner 和 JDBC Connector,这是一张按id更新和删除的 Upsert Sink。

MySQL DDL 中的:

PRIMARYKEY(id)

才是数据库真正执行的唯一约束。

两边都要有,而且 Key 语义必须一致。

七、第一次启动 Paimon 到 MySQL 的出仓作业

在用于运行同步作业的 SQL Client 会话中确认切回 Streaming 模式,再执行:

SET'execution.runtime-mode'='streaming';INSERTINTOmysql_ordersSELECTid,amount,statusFROMorders/*+ OPTIONS( 'scan.mode' = 'latest-full', 'consumer-id' = 'mysql-orders-v1' ) */;

这不是执行完马上退出的普通查询,而是一条持续运行的流作业。

latest-full可以先按下面这句话理解:

启动时读取最新 Snapshot 的当前全量,之后继续读取新产生的变化。

所以第一次启动时,前面已经存在的10011002也会进入 MySQL,不需要等它们再次发生变化。

如果 SQL Client 以 Attached 模式运行,这个终端可能会一直被作业占用。后续写入 Paimon 和执行 UPDATE、DELETE 时,可以再开一个 SQL Client 会话。

八、先验证启动时全量是否写进 MySQL

在 MySQL 中执行:

SELECTid,amount,statusFROMserving.orders_rtORDERBYid;

等待 JDBC Sink Flush 后,预期看到:

idamountstatus
100180.00CREATED
100250.00CREATED

如果暂时查不到,不要第一时间判断数据丢了。本文配置了:

'sink.buffer-flush.interval'='1s'

低流量时,记录可能先在 JDBC Sink Buffer 中等待下一次定时 Flush。除此之外,还要看作业是否处于 RUNNING、Checkpoint 是否正常、TaskManager 日志有没有连接错误。

九、再写一条新订单,验证增量新增

在另一个 Flink SQL Client 会话中执行:

INSERTINTOordersVALUES(1003,CAST(120.00ASDECIMAL(10,2)),'CREATED');

Paimon 会产生新的 Snapshot。持续运行的 Source 发现它以后,把新订单交给 JDBC Sink。

再查 MySQL:

SELECTid,amount,statusFROMserving.orders_rtWHEREid=1003;

预期结果:

idamountstatus
1003120.00CREATED

到这里,我们已经验证了两种读取:

作业启动前存在的 1001、1002 -> 启动全量 作业启动后新增的 1003 -> 持续增量

十、同一个主键再写一次,验证更新

现在让订单1001从 80 元变成 100 元,状态变成PAID

对 Paimon Deduplicate 主键表,可以再次写入完整的新值:

INSERTINTOordersVALUES(1001,CAST(100.00ASDECIMAL(10,2)),'PAID');

虽然 SQL 写的是INSERT INTO,但id=1001已经存在。对当前表状态来说,这是同一主键的新版本,不应该再多出第二个1001

查询 Paimon:

SELECTid,amount,statusFROMordersWHEREid=1001;

预期只有一行:

idamountstatus
1001100.00PAID

再查询 MySQL:

SELECTid,amount,statusFROMserving.orders_rtWHEREid=1001;

MySQL 也应该仍然只有一行,并且金额和状态都已经更新。

为什么 MySQL 没有插出两行

JDBC Sink 对主键表使用 MySQL UPSERT,核心语句类似:

INSERTINTOorders_rt(id,amount,status)VALUES(?,?,?)ONDUPLICATEKEYUPDATEid=VALUES(id),amount=VALUES(amount),status=VALUES(status);

第一次id=1001不存在,走 INSERT。

第二次id=1001已存在,MySQL 主键发生重复,走 UPDATE 分支。

所以“更新订单”不代表 Connector 一定发送普通的:

UPDATEorders_rtSET...WHEREid=...;

对 MySQL JDBC Upsert Sink,新增和更新通常共用INSERT ... ON DUPLICATE KEY UPDATE

能不能直接执行 UPDATE Paimon 表

Paimon 1.4.2 在 Flink 1.17 及以上支持对主键表执行 UPDATE,但 UPDATE 是 Batch 模式操作,而且不能修改主键。

可以在另一个用于批处理的 SQL Client 会话中执行:

SET'execution.runtime-mode'='batch';UPDATEordersSETamount=CAST(110.00ASDECIMAL(10,2)),status='PAID'WHEREid=1001;

这次 Batch DML 提交新的 Paimon Snapshot 后,原来持续运行的出仓作业仍会继续发现并同步变化。

十一、删除订单,验证 DELETE

Paimon 的DELETE FROM同样在 Batch 模式执行,并且只支持满足条件的主键表与 Merge Engine。

在批处理 SQL Client 会话中执行:

SET'execution.runtime-mode'='batch';DELETEFROMordersWHEREid=1001;

删除完成后,先查 Paimon 当前状态:

SELECTid,amount,statusFROMordersORDERBYid;

1001应该已经不存在。

持续运行的 Source 读到 Delete Changelog 后,JDBC Sink 最终会按主键执行类似:

DELETEFROMorders_rtWHEREid=?;

再查 MySQL:

SELECTid,amount,statusFROMserving.orders_rtWHEREid=1001;

预期返回 0 行。

十二、把订单1001的完整过程串起来

现在不看源码,只看业务状态:

时刻Paimon 操作Paimon 当前结果MySQL 当前结果
T1首次写入1001, 80, CREATED1001, 80, CREATED出仓后相同
T2再写1001, 100, PAID1001, 100, PAIDUPSERT 后相同
T3DELETE1001不存在DELETE 后不存在

Paimon 和 MySQL 之间传递的不是“每隔一段时间复制整张表”,而是一条持续运行的变化链路。

从最终状态看:

新增:MySQL 出现一行 更新:同一主键的值被覆盖,行数不增加 删除:MySQL 中对应主键消失

至于 T1 到 T2 之间究竟发出了+U,还是-D/+I;它们有没有在同一个 Buffer 中合并;什么时候调用executeBatch(),留到下一篇继续拆。

十三、两个主键为什么缺一不可

把两边的定义放在一起:

-- Flink JDBC Sink DDLPRIMARYKEY(id)NOTENFORCED-- MySQL 物理表 DDLPRIMARYKEY(id)

它们分别解决不同问题:

主键使用者作用
Flink DDL 主键Planner、JDBC Connector识别 Upsert Key,允许处理 UPDATE 和 DELETE
MySQL 物理主键MySQL阻止重复 Key,触发 UPSERT 的 UPDATE 分支

只有 Flink 主键,没有 MySQL 主键

Connector 仍可能生成 UPSERT SQL,但 MySQL 找不到重复键冲突。同一个id再来一次时,目标表可能出现重复行,幂等恢复也失去基础。

只有 MySQL 主键,没有 Flink 主键

Flink 会把 Sink 当成 Append 模式。上游查询一旦包含 UPDATE 或 DELETE,规划阶段通常就会拒绝,或者无法得到期望的更新语义。

两边都有主键,但字段不一致

如果 Paimon 用(tenant_id, order_id)标识订单,MySQL 却只用order_id,不同租户的相同订单号会互相覆盖。

因此主键检查不能只看“都有 PRIMARY KEY”,还要检查字段数量、顺序、类型和业务语义。

十四、latest-fullconsumer-id分别在做什么

本文 Source 使用:

'scan.mode'='latest-full','consumer-id'='mysql-orders-v1'

latest-full

它解决第一次启动从哪里读:

先读取最新 Snapshot 的完整当前状态 再持续读取后续新变化

这适合第一次为一张空的 MySQL 服务表建立镜像。

consumer-id

它给这一路长期消费一个稳定身份,Paimon 可以据此管理 Consumer 相关进度和 Snapshot 保留。

但要注意:

consumer-id不是 Flink Checkpoint,也不能单独保证作业故障后精确恢复到某一条记录。

Flink Source Split、读取位置和算子状态仍然需要 Checkpoint 或 Savepoint 保存。

十五、为什么刚写入后 MySQL 可能还查不到

本文配置:

'sink.buffer-flush.max-rows'='100','sink.buffer-flush.interval'='1s'

记录到达 JDBC Sink 后,不一定立即访问 MySQL。

它可能先进入当前 Sink 子任务的内存 Buffer,直到下面任一条件发生:

收到的记录达到 100 条 等待时间达到 1 秒 Flink 开始做 Checkpoint 作业正常关闭并清理尾批

所以低流量实验中,看到约 1 秒的可见延迟通常是正常现象。

如果长时间仍没有数据,再检查:

  • Flink 作业是不是 RUNNING;
  • Source 有没有读到新 Snapshot;
  • Sink 有没有持续报 JDBC 异常;
  • MySQL 是否存在锁等待或连接耗尽;
  • Checkpoint 是否频繁失败;
  • 查询的是不是正确的数据库和表。

十六、第一次跑最常见的八类错误

1. 找不到 JDBC Factory

典型信息包含:

Could not find any factory for identifier 'jdbc'

优先检查 Flink JDBC Connector JAR 是否存在、版本是否匹配、集群进程是否已经重启。

2. 找不到 MySQL Driver

典型信息包含:

No suitable driver ClassNotFoundException: com.mysql.cj.jdbc.Driver

优先检查 MySQL Connector/J 是否在真正执行作业的 TaskManager Classpath 中。

3. MySQL 拒绝连接

常见原因包括 Host、端口、账号、密码、授权来源 Host、防火墙和 TLS 配置不正确。

不要只在本机测试 MySQL Client,要从 TaskManager 所在网络环境验证可达性。

4. Flink 中注册了主键,MySQL 却出现重复数据

执行:

SHOWCREATETABLEserving.orders_rt;

确认 MySQL 物理表真的存在 PRIMARY KEY 或语义完全一致的 UNIQUE KEY。

5. 新增能写,更新或删除规划失败

检查 Flink JDBC Sink DDL 是否声明主键,以及上游查询经过投影、Join、聚合后是否仍保留可用的 Upsert Key。

6. DELETE 执行后 MySQL 仍有数据

先确认 Paimon 的 DELETE DML 是否成功提交新 Snapshot,再确认 Changelog Producer 是否产生删除变化,最后检查两边 Key 是否一致。

7. 金额或字符串写入失败

检查 Paimon、Flink JDBC DDL 和 MySQL 三边的数据类型。

例如:

Paimon DECIMAL(10, 2) Flink Sink DECIMAL(10, 2) MySQL DECIMAL(10, 2)

字段名字相同不代表类型一定兼容。精度、长度、NULL 约束和字符集都可能导致失败。

8. 作业正常,但目标端看起来延迟很大

先看sink.buffer-flush.interval,再看 Sink 是否背压、MySQL 是否慢、Checkpoint 是否长时间执行。

十七、跑通以后做一次最小验收

不要只看到一条 INSERT 成功就宣布链路完成。至少验证下面六项:

验收项期望结果
启动全量作业启动前的 Paimon 当前数据进入 MySQL
增量新增新 Key 出现在 MySQL
同 Key 更新MySQL 仍只有一行,字段变为新值
删除MySQL 对应 Key 消失
NULL 和边界值类型、长度、精度符合预期
Checkpoint能持续成功,不只是作业显示 RUNNING

再补一项非常实用的核对:

-- Paimon 侧SELECTCOUNT(*)FROMorders;-- MySQL 侧SELECTCOUNT(*)FROMserving.orders_rt;

行数相同只是第一步,不足以证明内容完全一致。正式上线还要对主键集合、关键字段、删除和抽样明细做对账。

十八、这一篇先记住五句话

1. mysql_orders 是 Flink 逻辑表,orders_rt 才是 MySQL 物理表 2. Flink DDL 主键决定 Upsert 语义,MySQL 主键真正执行唯一约束 3. latest-full 先读当前全量,再持续读后续变化 4. 新增和更新通常通过 MySQL UPSERT 落地,删除按主键执行 DELETE 5. 写入先进入 JDBC Sink Buffer,低流量时不一定立刻在 MySQL 可见

到这里,我们已经从零跑通了一条最小的 Paimon 到 MySQL 出仓链路。

下一篇不再停留在 DDL 层,而是拿四条真实变化逐步跟进源码:

Paimon Source 怎样一条条发 RowData GenericJdbcSinkFunction.invoke() 为什么每次只收一条 TableBufferReducedStatementExecutor 怎样按主键覆盖 为什么收到 100 条输入,最终不一定执行 100 条 MySQL DML addBatch() 和 executeBatch() 到底分别做了什么

把这段过程看清楚以后,“出仓到 MySQL 到底是一条一条还是一批一批”就不再只是一句结论,而是一条能从源码、日志和 MySQL 现象互相验证的完整链路。


本篇关键配置与资料位置

  • docs/content/flink/sql-write.md:Paimon 的 INSERT、UPDATE、DELETE DML 说明
  • JdbcDynamicTableSink.java:Flink JDBC Sink 的 ChangelogMode 和主键校验
  • MySqlDialect.java:生成 MySQLINSERT ... ON DUPLICATE KEY UPDATE
  • JdbcConnectorOptions.java:定义 Buffer Flush、时间间隔和重试配置
  • Flink 1.20 JDBC Connector 文档:依赖、Upsert 模式、主键要求和 Connector Options

本文基于 Apache Paimon 1.4.2、Apache Flink 1.20.1 和 Flink JDBC Connector 3.3.0-1.20。不同版本的依赖坐标、默认值和类名可能变化,实际部署时请以对应版本文档和源码为准。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询