2026/9/4 6:02:36

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

Apache Paimon 数据出仓源码导读(六):UPDATE 与 DELETE 如何落库:主键归并、UPSERT 与跨批空窗 上一篇我们跟着四条RowData走完了 JDBC Sink记录逐条进入按主键进入 Buffer满足条件后通过executeBatch()批量访问 MySQL。这一篇继续追问UPDATE 到了 MySQL为什么不一定执行普通 UPDATEDELETE 为什么只需要主键同一个订单先删后插有时只执行一次 UPSERT有时却真的先 DELETE、再 UPSERT答案藏在两个地方主键 Buffer 只保留同一 Key 的最后动作 Flush 边界决定哪些动作有机会在同一批里相遇一、先看订单1001的三种业务动作假设订单表当前是idamountstatus100180.00CREATED之后依次发生新增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 或 DELETESink DDL 必须声明主键原因很直接。没有 KeyConnector 不知道UPDATE 应该覆盖 MySQL 的哪一行 DELETE 应该删除 MySQL 的哪一行所以 Paimon 主键表出仓到 MySQL 时Flink JDBC DDL 通常必须声明与结果表一致的主键。四、有主键后Builder 选择 Upsert 路径JdbcOutputFormatBuilder.build()根据 Key 是否存在分支if(dmlOptions.getKeyFields().isPresent()dmlOptions.getKeyFields().get().length0){// upsert queryreturnnewJdbcOutputFormat(connectionProvider,executionOptions,()-createBufferReduceExecutor(...));}else{// append only queryreturnnewJdbcOutputFormat(connectionProvider,executionOptions,()-createSimpleBufferedExecutor(...));}两条路的区别路径Buffer可处理变化MySQL 行为Append普通批量 INSERT只适合 INSERT ONLY重放可能重复插入Upsert按主键归并I、U、-DUPSERT 或按 Key DELETE五、新增和更新为什么共用 UPSERTMySQL Dialect 为主键 Sink 生成类似INSERTINTOorders_rt(id,amount,status)VALUES(?,?,?)ONDUPLICATEKEYUPDATEidVALUES(id),amountVALUES(amount),statusVALUES(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 id1001 是否存在 存在 - UPDATE 不存在 - INSERT但它有两个问题。多一次数据库往返每条记录先查再写网络和 MySQL 查询开销都会明显增加。存在并发竞争两个 Writer 可能同时看到“不存在”然后都尝试 INSERT。MySQL UPSERT 把判断和写入交给数据库的唯一约束在一条 DML 语义中完成更适合主键同步。七、为什么 DELETE 只需要主键删除 SQL 类似DELETEFROMorders_rtWHEREid?它不需要旧金额和旧状态。即使上游删除记录只有-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 核心是MapPrimaryKey, LastChange因此同一 Key 的动作序列最终只剩最后一个同一 Flush 周期内的输入Buffer 最后动作MySQL 最终 DMLI - UUPSERT 新值1 次 UPSERTI - -DDELETE1 次 DELETE-D - IUPSERT 新值1 次 UPSERTU - -DDELETE1 次 DELETE-D - I - UUPSERT 最后新值1 次 UPSERT注意这张表只在“动作都进入同一个 Buffer”时成立。十、inputProducer 的-D/I为什么可能只剩 UPSERT假设上游把订单 1001 从 80 改成 100PaimoninputChangelog 中出现-D[1001,80] I[1001,100]如果两条记录在同一 Flush 周期到达收到 -D - 1001 DELETE 收到 I - 1001 UPSERT 100覆盖 DELETEFlush 时最终只有UPSERT 1001100MySQL 不需要真的经历“先没有 1001再重新出现 1001”。十一、跨过 Flush 边界后行为为什么完全不同现在让-D到来后立刻触发 FlushBatch A-D[1001,80] -------- Flush 边界 -------- Batch BI[1001,100]Batch A 已经清空 Buffer 并访问 MySQL。Batch B 不可能回头覆盖上一批动作。最终执行Batch A - DELETE id1001 Batch B - UPSERT id1001, amount100两个批次之间在线查询可能短暂看到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.EntryRowData,Tuple2Boolean,RowDataentry: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)NOTENFORCEDMySQL 物理表也应该有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 列是否允许 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-rows100,sink.buffer-flush.interval30s,sink.parallelism1快速连续发送-D/I观察 MySQL 是否始终保持新值。场景 B强制每条都 Flushsink.buffer-flush.max-rows1,sink.buffer-flush.interval0每条输入都会触发一次 Flush更容易看到DELETE 后暂时查不到 下一条 UPSERT 后重新出现验证时不要只看最终结果要持续高频查询目标 Key记录每次结果和时间戳。否则最后只看到1001100会错过中间空窗。二十、值班时怎样定位 UPDATE 或 DELETE 问题现象第一检查点第二检查点更新后出现重复行MySQL 真实 UNIQUE / PRIMARY KEYFlink Sink DDL KeyUPDATE 规划失败Sink 是否声明主键查询结果是否仍保留 Upsert KeyDELETE 没生效上游是否产生-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 BatchTableSimpleStatementExecutor.javaPreparedStatement 参数绑定和 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。