上一篇,我们跟着四条RowData走完了 JDBC Sink:记录逐条进入,按主键进入 Buffer,满足条件后通过executeBatch()批量访问 MySQL。
这一篇继续追问:
UPDATE 到了 MySQL,为什么不一定执行普通 UPDATE?DELETE 为什么只需要主键?同一个订单先删后插,有时只执行一次 UPSERT,有时却真的先 DELETE、再 UPSERT?
答案藏在两个地方:
主键 Buffer 只保留同一 Key 的最后动作 Flush 边界决定哪些动作有机会在同一批里相遇一、先看订单1001的三种业务动作
假设订单表当前是:
| id | amount | status |
|---|---|---|
| 1001 | 80.00 | CREATED |
之后依次发生:
新增:1001, 80.00, CREATED 更新:1001, 100.00, PAID 删除:1001MySQL 目标表要维护的是当前状态:
新增后 -> 有 1001 更新后 -> 仍只有一行 1001,但值变了 删除后 -> 没有 1001JDBC Sink 不需要在 MySQL 里永久保存每一张变化凭证,只需要把最终状态正确落下去。
二、Flink JDBC Sink 实际接收哪几种 RowKind
有主键的 JDBC Sink 在JdbcDynamicTableSink.getChangelogMode()中声明:
returnChangelogMode.newBuilder().addContainedKind(RowKind.INSERT).addContainedKind(RowKind.DELETE).addContainedKind(RowKind.UPDATE_AFTER).build();也就是:
+I INSERT +U UPDATE_AFTER -D DELETE这里没有-U UPDATE_BEFORE。
对一张按主键覆盖的 MySQL 服务表,JDBC Sink 真正需要知道的是:
这个 Key 的新值是什么? 或者,这个 Key 是否应该被删除?旧值-U对无主键聚合和撤回计算很重要,但对最终主键覆盖通常不是必需输入。
三、没有主键,UPDATE 和 DELETE 为什么进不去
源码会校验:
checkState(ChangelogMode.insertOnly().equals(requestedMode)||dmlOptions.getKeyFields().isPresent(),"please declare primary key for sink table when query contains update/delete record.");翻译成人话:
如果上游只有 INSERT,可以没有主键 如果上游包含 UPDATE 或 DELETE,Sink DDL 必须声明主键原因很直接。
没有 Key,Connector 不知道:
UPDATE 应该覆盖 MySQL 的哪一行 DELETE 应该删除 MySQL 的哪一行所以 Paimon 主键表出仓到 MySQL 时,Flink JDBC DDL 通常必须声明与结果表一致的主键。
四、有主键后,Builder 选择 Upsert 路径
JdbcOutputFormatBuilder.build()根据 Key 是否存在分支:
if(dmlOptions.getKeyFields().isPresent()&&dmlOptions.getKeyFields().get().length>0){// upsert queryreturnnewJdbcOutputFormat<>(connectionProvider,executionOptions,()->createBufferReduceExecutor(...));}else{// append only queryreturnnewJdbcOutputFormat<>(connectionProvider,executionOptions,()->createSimpleBufferedExecutor(...));}两条路的区别:
| 路径 | Buffer | 可处理变化 | MySQL 行为 |
|---|---|---|---|
| Append | 普通批量 INSERT | 只适合 INSERT ONLY | 重放可能重复插入 |
| Upsert | 按主键归并 | +I、+U、-D | UPSERT 或按 Key DELETE |
五、新增和更新为什么共用 UPSERT
MySQL Dialect 为主键 Sink 生成类似:
INSERTINTO`orders_rt`(`id`,`amount`,`status`)VALUES(?,?,?)ONDUPLICATEKEYUPDATE`id`=VALUES(`id`),`amount`=VALUES(`amount`),`status`=VALUES(`status`)逐段读:
INSERTINTO...VALUES(?,?,?)先尝试插入这份新状态。
ONDUPLICATEKEYUPDATE...如果 MySQL 的 PRIMARY KEY 或 UNIQUE KEY 已经存在,就把原行更新成本次参数值。
因此:
+I[1001,80,CREATED] -> UPSERT +U[1001,100,PAID] -> UPSERT两种 RowKind 最终可以使用同一条 PreparedStatement 模板。
六、为什么不先 SELECT,再决定 INSERT 或 UPDATE
看起来最直观的实现是:
SELECT id=1001 是否存在 存在 -> UPDATE 不存在 -> INSERT但它有两个问题。
多一次数据库往返
每条记录先查再写,网络和 MySQL 查询开销都会明显增加。
存在并发竞争
两个 Writer 可能同时看到“不存在”,然后都尝试 INSERT。
MySQL UPSERT 把判断和写入交给数据库的唯一约束,在一条 DML 语义中完成,更适合主键同步。
七、为什么 DELETE 只需要主键
删除 SQL 类似:
DELETEFROM`orders_rt`WHERE`id`=?它不需要旧金额和旧状态。
即使上游删除记录只有:
-D[1001]只要 Key 完整,JDBC Sink 就能生成准确的 WHERE 参数。
这也是复合主键特别需要小心的原因。
如果 MySQL 主键是:
(tenant_id, order_id)Delete Changelog 就必须能提取两个字段。少一个都无法唯一定位目标行。
八、RowKind 怎样变成“加入”或“撤回”
TableBufferReducedStatementExecutor用一个 Boolean 标记最终动作:
privatebooleanchangeFlag(RowKindrowKind){switch(rowKind){caseINSERT:caseUPDATE_AFTER:returntrue;caseDELETE:caseUPDATE_BEFORE:returnfalse;default:thrownewUnsupportedOperationException(...);}}可以简化为:
true -> 加入组 -> UPSERT false -> 撤回组 -> DELETE虽然内部执行器认识UPDATE_BEFORE,但标准 Table Sink 声明的目标 ChangelogMode 仍是+I/+U/-D。Planner 会尽量把上游变化调整为 Sink 能消费的形式。
九、同一 Buffer 内,后来的动作覆盖前面的动作
Buffer 核心是:
Map<PrimaryKey, LastChange>因此同一 Key 的动作序列,最终只剩最后一个:
| 同一 Flush 周期内的输入 | Buffer 最后动作 | MySQL 最终 DML |
|---|---|---|
+I -> +U | UPSERT 新值 | 1 次 UPSERT |
+I -> -D | DELETE | 1 次 DELETE |
-D -> +I | UPSERT 新值 | 1 次 UPSERT |
+U -> -D | DELETE | 1 次 DELETE |
-D -> +I -> +U | UPSERT 最后新值 | 1 次 UPSERT |
注意,这张表只在“动作都进入同一个 Buffer”时成立。
十、inputProducer 的-D/+I为什么可能只剩 UPSERT
假设上游把订单 1001 从 80 改成 100,PaimoninputChangelog 中出现:
-D[1001,80] +I[1001,100]如果两条记录在同一 Flush 周期到达:
收到 -D -> 1001 = DELETE 收到 +I -> 1001 = UPSERT 100,覆盖 DELETEFlush 时最终只有:
UPSERT 1001=100MySQL 不需要真的经历“先没有 1001,再重新出现 1001”。
十一、跨过 Flush 边界后,行为为什么完全不同
现在让-D到来后立刻触发 Flush:
Batch A:-D[1001,80] -------- Flush 边界 -------- Batch B:+I[1001,100]Batch A 已经清空 Buffer 并访问 MySQL。Batch B 不可能回头覆盖上一批动作。
最终执行:
Batch A -> DELETE id=1001 Batch B -> UPSERT id=1001, amount=100两个批次之间,在线查询可能短暂看到:
1001 不存在随后才重新出现新值。
这就是“最终状态正确”和“中间过程无空窗”之间的区别。
十二、哪些事件会把-D/+I切到两个批次
Flush 边界可能来自:
sink.buffer-flush.max-rows恰好达到阈值;sink.buffer-flush.interval定时器恰好触发;- Checkpoint 到来并强制 Flush;
- 作业结束触发 Close Flush;
- 不同记录被路由到不同处理阶段或作业。
因此不能只看:
-D 和 +I 在 Changelog 中是否相邻还要看它们进入 JDBC Sink 时有没有跨过实际 Flush 边界。
十三、UPSERT Batch 和 DELETE Batch 的执行顺序
执行器源码是:
for(Map.Entry<RowData,Tuple2<Boolean,RowData>>entry:reduceBuffer.entrySet()){if(entry.getValue().f0){upsertExecutor.addToBatch(entry.getValue().f1);}else{deleteExecutor.addToBatch(entry.getKey());}}upsertExecutor.executeBatch();deleteExecutor.executeBatch();reduceBuffer.clear();也就是一次 Flush 内:
先执行 UPSERT Batch 再执行 DELETE Batch同一个 Key 只保留最后动作,所以不会既出现在 UPSERT 组又出现在 DELETE 组。
不同 Key 则可能分处两组,例如:
1001 -> UPSERT 1002 -> DELETE这两组并不是一笔与 Flink Checkpoint 原子绑定的数据库事务。UPSERT 组成功、DELETE 组失败时,重试可能再次执行部分动作。故障恢复篇会继续展开。
十四、不同 Changelog Producer 对 MySQL 的影响
| Producer | 常见更新输出 | JDBC Sink 需要的最终动作 | 主要关注点 |
|---|---|---|---|
none | 新值 Upsert Changelog | UPSERT 新值 | 适合只维护目标当前状态 |
input | 上游可能是-D/+I | 同批可归并;跨批会先删后插 | 可能产生短暂空窗 |
lookup | 常见-U/+U | Planner 保留新值+U | 旧值生成有额外成本 |
full-compaction | 延迟产生净变化 | UPSERT / DELETE | MySQL 可见延迟受 Full Compaction 影响 |
如果 MySQL 目标只关心每个 Key 的当前状态,通常不需要为了 JDBC Sink 强行生成旧值。
但如果下游还有聚合、审计或撤回计算,就不能只从 MySQL Sink 的需求选择 Producer。
十五、复合主键必须三边完全一致
假设 Paimon 订单唯一键是:
(tenant_id, order_id)那么 Flink JDBC DDL 应该声明:
PRIMARYKEY(tenant_id,order_id)NOTENFORCEDMySQL 物理表也应该有:
PRIMARYKEY(tenant_id,order_id)常见错误:
| 错误 | 后果 |
|---|---|
MySQL 只用order_id | 不同租户相互覆盖 |
Flink Sink 只声明order_id | Buffer 先把不同租户错误归并 |
| MySQL Key 多一个 Sink 没有的字段 | UPSERT 和 DELETE 无法按同一语义定位 |
| 字符串排序或大小写规则不同 | 逻辑相同的 Key 在两端可能判断不同 |
十六、NULL、类型和字段顺序也会影响落库
主键字段通常不允许 NULL,但非主键字段仍要处理:
Paimon STRING -> MySQL VARCHAR 长度是否足够 Paimon DECIMAL(10,2) -> MySQL 精度是否一致 Paimon TIMESTAMP -> MySQL 时区和精度是否一致 Paimon NULL -> MySQL 列是否允许 NULLFlink 逻辑表字段顺序要与 Connector 生成参数的字段映射一致。
生产上线前至少验证:
- 最大字符串长度;
- 最大和最小金额;
- NULL 值;
- 多字节字符;
- 时间边界和时区;
- 复合主键所有字段。
十七、结果幂等不代表所有副作用幂等
重复 UPSERT 相同参数,主表最终通常相同。
重复 DELETE 同一 Key,主表最终也都是不存在。
但如果目标表有:
- INSERT / UPDATE / DELETE Trigger;
- 审计流水;
- 版本号自增;
- 更新时间强制刷新;
- 级联删除;
- 外部通知;
重复执行可能产生额外副作用。
所以“主表最终值正确”不能替代对 Trigger、审计表和下游通知的检查。
十八、多个 Writer 同时改一个 Key 会怎样
假设 Flink 出仓作业和业务服务都能写orders_rt。
时间线:
1. Flink 已写 PAID 2. 业务服务改成 REFUNDED 3. Flink 故障恢复,重放旧的 PAID 4. UPSERT 把 REFUNDED 覆盖回 PAID对“同一条 Flink 记录重复执行”,UPSERT 是幂等的。
对“多个系统竞争修改”,它只是最后写入者覆盖前者,不会自动判断哪个版本更新。
如果目标表存在多 Writer,需要增加:
- 版本号或事件时间条件;
- 单 Writer 所有权;
- 冲突检测;
- 独立影子表;
- 明确的数据权威来源。
十九、怎样做一个跨批空窗实验
为了复现,不要依赖运气,可以人为缩小 Flush 条件。
场景 A:尽量让-D/+I同批
'sink.buffer-flush.max-rows'='100','sink.buffer-flush.interval'='30s','sink.parallelism'='1'快速连续发送-D/+I,观察 MySQL 是否始终保持新值。
场景 B:强制每条都 Flush
'sink.buffer-flush.max-rows'='1','sink.buffer-flush.interval'='0'每条输入都会触发一次 Flush,更容易看到:
DELETE 后暂时查不到 下一条 UPSERT 后重新出现验证时不要只看最终结果
要持续高频查询目标 Key,记录每次结果和时间戳。否则最后只看到1001=100,会错过中间空窗。
二十、值班时怎样定位 UPDATE 或 DELETE 问题
| 现象 | 第一检查点 | 第二检查点 |
|---|---|---|
| 更新后出现重复行 | MySQL 真实 UNIQUE / PRIMARY KEY | Flink Sink DDL Key |
| UPDATE 规划失败 | Sink 是否声明主键 | 查询结果是否仍保留 Upsert Key |
| DELETE 没生效 | 上游是否产生-D | Key 字段是否完整一致 |
| 偶尔先消失后出现 | 是否为input的-D/+I | 是否跨过 Flush 边界 |
| 主表正确但审计重复 | Trigger / 审计逻辑 | 是否发生重试或恢复重放 |
| 其他业务修改被覆盖 | 是否存在多 Writer | 是否需要版本冲突控制 |
二十一、最后记住三条边界
第一,+I 和 +U 都可以走 MySQL UPSERT,-D 按主键 DELETE 第二,同一 Buffer 内同 Key 只保留最后动作;跨过 Flush 后无法再合并 第三,最终状态正确不代表中间没有空窗,也不代表所有副作用都幂等下一篇,我们故意让任务在最麻烦的位置失败:
JDBC Batch 已经写进 MySQL 新的 Checkpoint 却还没有全局成功然后观察为什么 Source 会重放、为什么普通 JDBC Sink 是 At-Least-Once,以及 UPSERT / DELETE 到底在什么条件下能把重复执行收敛。
本篇关键源码位置
JdbcDynamicTableSink.java:声明+I/+U/-D并校验更新流主键JdbcOutputFormatBuilder.java:根据 Key 选择 Upsert 或 Append 路径TableBufferReducedStatementExecutor.java:按 Key 保存最后动作并拆分 UPSERT / DELETE BatchTableSimpleStatementExecutor.java:PreparedStatement 参数绑定和 Batch 执行MySqlDialect.java:生成INSERT ... ON DUPLICATE KEY UPDATEAbstractDialect.java:生成按主键 DELETE SQL
本文基于 Apache Paimon 1.4.2、Apache Flink 1.20.1 和 Flink JDBC Connector 3.3.0-1.20。不同数据库 Dialect 的 UPSERT 语法不同,本文 MySQL 结论不能直接套用到 PostgreSQL、Oracle 或 SQL Server。