上一篇,我们从源码角度总览了 Paimon 出仓到 MySQL 时的 UPSERT、DELETE、批量写入和故障重放。
不过,如果还没有亲手跑过这条链路,一上来就看RowKind、JdbcOutputFormat、executeBatch()和 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这里有三张“表”,但它们并不都是同一种东西:
| 名字 | 存在哪里 | 它是什么 |
|---|---|---|
orders | Paimon | 真正保存订单数据的湖表 |
mysql_orders | Flink Catalog | 指向 MySQL 的 JDBC 表定义,本身不保存另一份数据 |
orders_rt | MySQL | 真正接收出仓结果的物理表 |
很多初学者会误以为执行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 ENFORCED | Flink 不负责运行时检查唯一性 |
path | Paimon 表在文件系统中的存储位置 |
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;预期当前结果是:
| id | amount | status |
|---|---|---|
| 1001 | 80.00 | CREATED |
| 1002 | 50.00 | CREATED |
这里看到的是 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 的当前全量,之后继续读取新产生的变化。所以第一次启动时,前面已经存在的1001和1002也会进入 MySQL,不需要等它们再次发生变化。
如果 SQL Client 以 Attached 模式运行,这个终端可能会一直被作业占用。后续写入 Paimon 和执行 UPDATE、DELETE 时,可以再开一个 SQL Client 会话。
八、先验证启动时全量是否写进 MySQL
在 MySQL 中执行:
SELECTid,amount,statusFROMserving.orders_rtORDERBYid;等待 JDBC Sink Flush 后,预期看到:
| id | amount | status |
|---|---|---|
| 1001 | 80.00 | CREATED |
| 1002 | 50.00 | CREATED |
如果暂时查不到,不要第一时间判断数据丢了。本文配置了:
'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;预期结果:
| id | amount | status |
|---|---|---|
| 1003 | 120.00 | CREATED |
到这里,我们已经验证了两种读取:
作业启动前存在的 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;预期只有一行:
| id | amount | status |
|---|---|---|
| 1001 | 100.00 | PAID |
再查询 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, CREATED | 1001, 80, CREATED | 出仓后相同 |
| T2 | 再写1001, 100, PAID | 1001, 100, PAID | UPSERT 后相同 |
| T3 | DELETE1001 | 不存在 | 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-full和consumer-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 UPDATEJdbcConnectorOptions.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。不同版本的依赖坐标、默认值和类名可能变化,实际部署时请以对应版本文档和源码为准。