上一篇,我们从零搭好了 Paimon 到 MySQL 的最小链路,并验证了新增、更新和删除。
现在回到最容易让人困惑的问题:
Paimon 有十万条变化要写 MySQL,Flink 到底是一条一条写,还是一批一批写?
这个问题之所以争论不清,不是因为某一边一定错了,而是因为大家说的“写”不在同一层。
这一篇不再只给结论。我们把批大小暂时改成 4,跟着四条订单变化从 Paimon Source 出发,经过invoke()、batchCount、主键 Buffer、addBatch(),最后走到executeBatch()。
一、先给出不会混淆的四层答案
把链路分成四层:
| 层次 | 实际动作 | 逐条还是批量 |
|---|---|---|
| Paimon Source | 从 Split 中读出一条RowData,交给 Flink 下游 | 逐条 |
| Flink Sink Function | 每来一条调用一次invoke() | 逐条 |
| Sink 子任务 Buffer | 在多次invoke()之间积累,并按主键保留最后动作 | 攒批 |
| JDBC / MySQL | 多次addBatch()后执行executeBatch() | 批量 |
所以最准确的一句话是:
记录逐条进入 Sink,每个 Sink 子任务独立攒批,满足 Flush 条件后批量访问 MySQL。
二、先把实验参数缩小,肉眼观察一批数据
生产配置经常是:
'sink.buffer-flush.max-rows'='100','sink.buffer-flush.interval'='1s'但 100 条不适合手工跟踪。为了讲清楚原理,本文临时改成:
'sink.buffer-flush.max-rows'='4','sink.buffer-flush.interval'='10s','sink.parallelism'='1'含义是:
最多收到 4 条输入就 Flush 低流量时最多等 10 秒也会 Flush 只有 1 个 Sink 子任务,避免多个 Buffer 干扰观察实验值只是为了看清过程。不要把4、10s、并行度1原样当作生产最佳实践。
三、今天跟踪的四条变化
假设 Sink 在 10 秒内依次收到:
第 1 条:+I[1001, 80.00, CREATED] 第 2 条:+U[1001, 100.00, PAID] 第 3 条:+I[1002, 50.00, CREATED] 第 4 条:-D[1002, 50.00, CREATED]这里故意让1001和1002各出现两次。
这样可以同时观察两个数字:
输入记录数:4 不同主键数:2后面会看到,max-rows=4统计的是前者,而主键 Buffer 最终保存的是后者。
四、Paimon Source 不是一次吐出一个 List
流式读取 Paimon 时,Enumerator 先发现需要读取的 Snapshot,再把 Snapshot 规划成 Split。
可以先用物流类比理解:
Snapshot:这一批仓库变化的版本 Split:分给某个读取任务的一箱货 RowData:箱子里一件一件拿出的商品对应源码主线是:
DataTableStreamScan 找到下一个 Snapshot 并生成 Plan ContinuousFileSplitEnumerator 把 Plan 转成 FileStoreSourceSplit 并分配给 Reader FileStoreSourceReader 消费 Split,把其中的 RowData 逐条发给下游FileStoreSourceReader继承 Flink 的 SourceReader 基础实现。它的工作单位虽然是 Split,但交给下游的元素类型是:
RowData因此,“一个 Split 里有一万条变化”不等于下游invoke()一次收到一万个对象。
Split 是读取任务分工单位,RowData 才是算子之间传递的数据单位。
五、一条 RowData 里除了字段值,还有 RowKind
业务上看到的订单是:
id=1001, amount=100.00, status=PAIDFlink 内部的RowData还带着动作类型:
+I INSERT -U UPDATE_BEFORE +U UPDATE_AFTER -D DELETEJDBC Sink 不只要知道字段值,还要知道这条记录应该加入 UPSERT 组还是 DELETE 组。
例如:
+U[1001, 100.00, PAID]表示“把 1001 的当前值改成这一份新值”。
而:
-D[1002, 50.00, CREATED]表示“按主键删除 1002”。
六、Flink 每来一条,就调用一次invoke()
JDBC Table Sink 最终构造GenericJdbcSinkFunction。
它的invoke()很短:
publicvoidinvoke(Tvalue,Contextcontext)throwsIOException{outputFormat.writeRecord(value);}关键点不是代码长短,而是方法参数:
Tvalue这里没有:
List<T>values也就是说:
第 1 条 RowData -> invoke() 第 1 次 第 2 条 RowData -> invoke() 第 2 次 第 3 条 RowData -> invoke() 第 3 次 第 4 条 RowData -> invoke() 第 4 次因此,从 Flink 算子调用角度说“逐条进入 Sink”完全正确。
七、writeRecord()做了三件事
invoke()随后调用JdbcOutputFormat.writeRecord():
publicfinalsynchronizedvoidwriteRecord(Inrecord)throwsIOException{checkFlushException();InrecordCopy=copyIfNecessary(record);addToBatch(record,getExtractor().apply(recordCopy));batchCount++;if(executionOptions.getBatchSize()>0&&batchCount>=executionOptions.getBatchSize()){flush();}}先不被类名吓住,把它翻译成人话:
1. 检查之前的异步 Flush 有没有失败 2. 把当前这一条记录加入内部执行器 3. batchCount 加 1,达到阈值就 Flush每调用一次writeRecord(),batchCount就加一次。
但“加入内部执行器”不代表已经访问 MySQL。主键表还会先进入一张内存 Map。
八、为什么源码要复制一份 RowData
代码里有:
InrecordCopy=copyIfNecessary(record);初学者容易觉得多此一举。
Flink 为了减少对象创建,可以开启 Object Reuse。上游下一条记录到来时,可能复用并改写同一个 Java 对象。
如果 JDBC Sink 把原对象直接放进 Buffer:
刚放进去的是 1001 上游复用对象后把它改成 1002 Buffer 里原本的 1001 也可能被一起“变掉”所以需要根据序列化器和 Object Reuse 配置保留安全副本。
这是实现细节,却解释了为什么“Buffer 里保存一条记录”不只是把引用塞进集合那么简单。
九、主键模式使用 Map,不是简单 List
JDBC Sink DDL 声明主键后,Builder 会选择:
TableBufferReducedStatementExecutor它的核心 Buffer 是:
privatefinalMap<RowData,Tuple2<Boolean,RowData>>reduceBuffer=newHashMap<>();可以简化成:
Map<主键, 最后动作>每条数据到来时:
RowDatakey=keyExtractor.apply(record);booleanflag=changeFlag(record.getRowKind());reduceBuffer.put(key,Tuple2.of(flag,record));Map.put()的特性是:同一个 Key 再次写入,新值覆盖旧值。
所以 JDBC Sink 不会在同一个 Flush 周期内,为同一个订单无限堆积所有历史动作。
十、四条记录进入 Buffer 时发生了什么
逐条展开:
| 到达顺序 | 输入 | batchCount | reduceBuffer |
|---|---|---|---|
| 1 | +I[1001,80] | 1 | 1001 -> UPSERT 80 |
| 2 | +U[1001,100] | 2 | 1001 -> UPSERT 100 |
| 3 | +I[1002,50] | 3 | 1001 -> UPSERT 100;1002 -> UPSERT 50 |
| 4 | -D[1002] | 4 | 1001 -> UPSERT 100;1002 -> DELETE |
第四条到来后,同时出现:
batchCount = 4 reduceBuffer.size() = 2这两个数字不相等完全正常。
batchCount回答:
这个 Flush 周期收到了多少条输入?reduceBuffer.size()回答:
按主键归并后,还剩多少个最终动作?十一、max-rows=4为什么可能只执行两组参数
第四条让batchCount达到 4,于是触发:
flush();执行器遍历 Map:
if(entry.getValue().f0){upsertExecutor.addToBatch(entry.getValue().f1);}else{deleteExecutor.addToBatch(entry.getKey());}最终形成:
UPSERT Batch:1001 -> 100 DELETE Batch:1002所以一次实验里出现三个数:
| 指标 | 数值 |
|---|---|
| Sink 输入记录数 | 4 |
| 最终不同 Key 数 | 2 |
| Flush 次数 | 1 |
这就是“配置 100 条一批,MySQL 为什么不一定看到 100 个不同 Key”的根本原因。
十二、Flush 不只由最大行数触发
JDBC Sink 至少有四种 Flush 时机:
| 触发器 | 对应位置 | 作用 |
|---|---|---|
| 最大输入行数 | writeRecord() | 高流量时控制单批上限 |
| 时间间隔 | JdbcOutputFormat.open()启动的定时线程 | 低流量时控制可见延迟 |
| Checkpoint | snapshotState() | 保存状态前清空内存 Buffer |
| Close | close() | 有界任务或正常结束时处理尾批 |
所以真实批次大小通常小于等于max-rows。
例如配置 100 条,但只来了 7 条:
1 秒时间间隔先到 -> 7 条也会 Flush又例如来了 40 条时 Checkpoint 到达:
Checkpoint 先到 -> 40 条也会 Flushmax-rows是触发上限,不是每一批必须凑满的承诺。
十三、addBatch()和executeBatch()分别做什么
归并完成后,真正操作 PreparedStatement 的是TableSimpleStatementExecutor:
publicvoidaddToBatch(RowDatarecord)throwsSQLException{converter.toExternal(record,st);st.addBatch();}publicvoidexecuteBatch()throwsSQLException{st.executeBatch();}可以把 PreparedStatement 想成一张可重复填写的模板:
INSERTINTOorders_rt(id,amount,status)VALUES(?,?,?)ONDUPLICATEKEYUPDATE...对每一条最终动作:
设置第 1 组参数 -> addBatch() 设置第 2 组参数 -> addBatch() 设置第 3 组参数 -> addBatch() ... 最后 -> executeBatch()addBatch()只是把当前参数加入 JDBC Batch,不等于数据库已经执行。
executeBatch()才把累积的参数交给 Driver 执行。
十四、一次 Flush 内部还有两个 Batch
主键表需要同时处理新增、更新和删除。
因此TableBufferReducedStatementExecutor.executeBatch()会拆成:
upsertExecutor deleteExecutor源码顺序是:
upsertExecutor.executeBatch();deleteExecutor.executeBatch();reduceBuffer.clear();所以“一次 Flush”不代表只能调用一次数据库 PreparedStatement。
更准确地说:
一次 Flush ├─ 执行当前 UPSERT Batch └─ 执行当前 DELETE Batch同一个 Key 在 Map 中只剩最后动作,因此不会同时进入两组。
十五、executeBatch()不一定等于一条超长 SQL
Flink Connector 能确定的是:
PreparedStatement 设置多组参数 多次 addBatch() 调用 executeBatch()Driver 具体怎样通过 MySQL 协议发送,还要看 Driver 实现和连接参数。
Flink JDBC Connector 3.3.0 的 MySQL Dialect 会在 URL 未显式配置时追加:
rewriteBatchedStatements=true这允许 Connector/J 在适用场景下重写批处理,减少网络往返。
但不要把它理解成:
Flink 一批 100 条 = MySQL 慢日志必然出现一条带 100 个 VALUES 的 SQL语句类型、Driver 版本和重写条件都会影响最终协议形态。
排查时要区分:
Flink 输入了多少条 Buffer 归并后有多少 Key JDBC executeBatch 调了几次 Driver 发了多少网络请求 MySQL 最终影响多少行它们不是同一个指标。
十六、并行度大于 1 时,不是大家共用一个 Buffer
假设:
'sink.parallelism'='3','sink.buffer-flush.max-rows'='100'实际结构是:
每个 Sink 子任务通常有自己的:
JdbcOutputFormat reduceBuffer batchCount 定时 Flush 线程 JDBC 连接不是整个作业累计 100 条就统一 Flush,而是:
子任务 0 收到 100 条 -> 子任务 0 Flush 子任务 1 只收到 12 条 -> 继续等待时间或 Checkpoint 子任务 2 收到 100 条 -> 子任务 2 Flush因此提高并行度同时会提高潜在 MySQL 连接数和并发 Batch 数量。
十七、为什么热点 Key 会让“批大小”看起来很奇怪
假设一个子任务收到 100 条输入,但全都在更新订单1001:
batchCount = 100 reduceBuffer.size() = 1达到阈值后仍会 Flush,但真正只剩1001的最后值。
反过来,如果 100 条都是不同主键:
batchCount = 100 reduceBuffer.size() = 100这两种流量的输入条数一样,对 MySQL 的实际写入压力却完全不同。
所以压测数据不能只有“每秒多少条”,还要描述:
- 不同 Key 比例;
- 同 Key 更新频率;
- 更新与删除占比;
- Key 是否倾斜;
- 每行字段大小。
十八、怎样验证当前到底是逐条还是批量
可以按下面顺序验证。
第一步:把并行度固定为 1
先排除多个子任务各自 Flush 造成的干扰。
第二步:把max-rows改为 4,时间间隔改大
'sink.buffer-flush.max-rows'='4','sink.buffer-flush.interval'='30s'让行数更容易成为第一个触发器。
第三步:连续写入四个不同 Key
先避免同 Key 归并,让最终动作数更直观。
第四步:再连续更新同一个 Key 四次
对比batchCount=4,但最终只剩一个主键动作的情况。
第五步:观察多个层次
不要只看 MySQL 慢日志。至少同时观察:
Flink Sink 输入速率 Checkpoint 和 Flush 时长 MySQL 执行次数与影响行数 JDBC 连接数 网络往返与锁等待十九、最常见的六个误区
| 误区 | 正确认识 |
|---|---|
| Source 一次规划一个 Split,所以 Sink 一次收到一个 Split | Split 在 Source 内部逐条读成 RowData |
invoke()逐条调用,所以 MySQL 一定逐条访问 | 多次 invoke 之间会积累 Buffer |
max-rows=100就一定有 100 个不同 Key | 统计的是输入条数,同 Key 会归并 |
| 一次 Flush 就只有一条 SQL | UPSERT 和 DELETE 使用不同 PreparedStatement Batch |
executeBatch()就一定是一条多 VALUES SQL | 最终协议形态由 Driver 决定 |
| 并行度 4 仍只有一份 Buffer | 每个 Sink 子任务各有 Buffer、计数和连接 |
二十、最后把源码调用链连起来
一条订单变化的调用主线可以压缩成:
DataTableStreamScan 找到需要消费的 Snapshot ↓ ContinuousFileSplitEnumerator 规划并分配 FileStoreSourceSplit ↓ FileStoreSourceReader 从 Split 逐条发出 RowData ↓ GenericJdbcSinkFunction.invoke(row) 每条记录调用一次 ↓ JdbcOutputFormat.writeRecord(row) 加入执行器,batchCount++ ↓ TableBufferReducedStatementExecutor.addToBatch(row) 按主键放进 Map,后来的覆盖前面的 ↓ JdbcOutputFormat.flush() 行数 / 时间 / Checkpoint / Close 触发 ↓ TableSimpleStatementExecutor PreparedStatement.addBatch() PreparedStatement.executeBatch() ↓ MySQL现在再回答开头的问题:
Paimon Source:逐条发 Flink invoke:逐条收 Sink Buffer:按主键攒批 JDBC / MySQL:Flush 时批量执行“逐条”和“批量”都对,关键是必须说明自己讨论的是哪一层。
下一篇,我们继续深入这张主键 Map:
+I、+U 为什么都进 UPSERT -D 为什么只需要主键 同一 Buffer 内的 -D/+I 为什么可能只剩一次 UPSERT 跨过 Flush 边界后为什么可能先 DELETE、后 UPSERT 两个批次之间为什么会出现短暂空窗本篇关键源码位置
DataTableStreamScan.java:持续发现和规划下一个 Paimon SnapshotContinuousFileSplitEnumerator.java:生成并分配 FileStoreSourceSplitFileStoreSourceReader.java:消费 Split,逐条向下游发出 RowDataGenericJdbcSinkFunction.java:每条记录调用writeRecord(),Checkpoint 时调用flush()JdbcOutputFormat.java:维护batchCount、定时任务、行数 Flush、重试和 Close FlushTableBufferReducedStatementExecutor.java:使用 Map 按主键归并最后动作TableSimpleStatementExecutor.java:设置 PreparedStatement 参数并调用addBatch()、executeBatch()MySqlDialect.java:补充rewriteBatchedStatements=true
本文 Paimon 源码基于 Apache Paimon 1.4.2;JDBC 部分基于 Flink JDBC Connector 3.3.0-1.20 的 v3.3.0 源码。升级版本时请重新核对类名、默认值和 Driver 行为。