news 2026/9/4 6:02:25

Apache Paimon 数据出仓源码导读(六):UPDATE 与 DELETE 如何落库:主键归并、UPSERT 与跨批空窗

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Paimon 数据出仓源码导读(六):UPDATE 与 DELETE 如何落库:主键归并、UPSERT 与跨批空窗

上一篇,我们跟着四条RowData走完了 JDBC Sink:记录逐条进入,按主键进入 Buffer,满足条件后通过executeBatch()批量访问 MySQL。

这一篇继续追问:

UPDATE 到了 MySQL,为什么不一定执行普通 UPDATE?DELETE 为什么只需要主键?同一个订单先删后插,有时只执行一次 UPSERT,有时却真的先 DELETE、再 UPSERT?

答案藏在两个地方:

主键 Buffer 只保留同一 Key 的最后动作 Flush 边界决定哪些动作有机会在同一批里相遇

一、先看订单1001的三种业务动作

假设订单表当前是:

idamountstatus
100180.00CREATED

之后依次发生:

新增:1001, 80.00, CREATED 更新:1001, 100.00, PAID 删除:1001

MySQL 目标表要维护的是当前状态:

新增后 -> 有 1001 更新后 -> 仍只有一行 1001,但值变了 删除后 -> 没有 1001

JDBC 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-DUPSERT 或按 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 -> +UUPSERT 新值1 次 UPSERT
+I -> -DDELETE1 次 DELETE
-D -> +IUPSERT 新值1 次 UPSERT
+U -> -DDELETE1 次 DELETE
-D -> +I -> +UUPSERT 最后新值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,覆盖 DELETE

Flush 时最终只有:

UPSERT 1001=100

MySQL 不需要真的经历“先没有 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 ChangelogUPSERT 新值适合只维护目标当前状态
input上游可能是-D/+I同批可归并;跨批会先删后插可能产生短暂空窗
lookup常见-U/+UPlanner 保留新值+U旧值生成有额外成本
full-compaction延迟产生净变化UPSERT / DELETEMySQL 可见延迟受 Full Compaction 影响

如果 MySQL 目标只关心每个 Key 的当前状态,通常不需要为了 JDBC Sink 强行生成旧值。

但如果下游还有聚合、审计或撤回计算,就不能只从 MySQL Sink 的需求选择 Producer。

十五、复合主键必须三边完全一致

假设 Paimon 订单唯一键是:

(tenant_id, order_id)

那么 Flink JDBC DDL 应该声明:

PRIMARYKEY(tenant_id,order_id)NOTENFORCED

MySQL 物理表也应该有:

PRIMARYKEY(tenant_id,order_id)

常见错误:

错误后果
MySQL 只用order_id不同租户相互覆盖
Flink Sink 只声明order_idBuffer 先把不同租户错误归并
MySQL Key 多一个 Sink 没有的字段UPSERT 和 DELETE 无法按同一语义定位
字符串排序或大小写规则不同逻辑相同的 Key 在两端可能判断不同

十六、NULL、类型和字段顺序也会影响落库

主键字段通常不允许 NULL,但非主键字段仍要处理:

Paimon STRING -> MySQL VARCHAR 长度是否足够 Paimon DECIMAL(10,2) -> MySQL 精度是否一致 Paimon TIMESTAMP -> MySQL 时区和精度是否一致 Paimon NULL -> MySQL 列是否允许 NULL

Flink 逻辑表字段顺序要与 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 KEYFlink Sink DDL Key
UPDATE 规划失败Sink 是否声明主键查询结果是否仍保留 Upsert Key
DELETE 没生效上游是否产生-DKey 字段是否完整一致
偶尔先消失后出现是否为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 Batch
  • TableSimpleStatementExecutor.java:PreparedStatement 参数绑定和 Batch 执行
  • MySqlDialect.java:生成INSERT ... ON DUPLICATE KEY UPDATE
  • AbstractDialect.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。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/4 6:02:13

从零配置代码生成工具链:环境搭建、认证与验证全流程

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/4 6:02:10

从代码重构到架构优化:如何识别并重构软件中的设计债务

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/4 6:01:24

YOLOv5批量推理优化:从单图到多图并行的吞吐量提升实战

一张图12ms,100张图却要5秒?你的GPU可能在“摸鱼”!本文基于2026年最新实测数据,深入剖析YOLOv5批量推理的优化全链路——从动态批处理到TensorRT深度调优,从多GPU并行到服务化部署,手把手带你将推理吞吐量提升5-10倍。 一、问题的真相:你的GPU利用率为什么只有35%? 2…

作者头像 李华
网站建设 2026/9/4 6:00:13

零基础部署 OpenClaw,图形化安装避开 Python/Node 环境坑

OpenClaw 一键安装包&#xff5c;图形化一键部署&#xff0c;告别复杂环境配置 适配系统&#xff1a;Windows10/11 64 位、macOS 12 当前版本&#xff1a;Windows v3.1.0&#xff5c;macOS v2.7.9 OpenClaw 提供图形化一键部署方案&#xff0c;全程可视化交互&#xff0c;不需要…

作者头像 李华
网站建设 2026/9/4 5:59:51

从光流到RIFE:视频补帧技术原理与动作视频实战指南

“补帧”到底解决了什么问题&#xff1f;先问一个很多视频爱好者和 UP 主都遇到过的问题&#xff1a;你下载到一段 30fps 的舞蹈视频&#xff0c;画面里人物动作很快&#xff0c;比如劈叉、旋转、踢腿&#xff0c;逐帧看的时候总觉得跳跃感明显&#xff1b;你想做慢动作回放&am…

作者头像 李华
网站建设 2026/9/4 5:58:24

VtorShell变量与流程控制实战:从单条命令到自动化脚本

VtorShell 这类脚本解释器版本里&#xff0c;最值得关注的改动就是“支持变量与流程控制”。它让脚本从“固定写死的一条命令”变成“可以根据当前状态决定下一步跑什么”&#xff0c;适合正在做自动化脚本、批量任务或想把一次性命令整理成可复用脚本的人。只加变量不算难&…

作者头像 李华