news 2026/9/4 6:03:02

Apache Paimon 数据出仓源码导读(五):一条数据怎样写进 MySQL:RowData、Sink 子任务与 JDBC Batch

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Paimon 数据出仓源码导读(五):一条数据怎样写进 MySQL:RowData、Sink 子任务与 JDBC Batch

上一篇,我们从零搭好了 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 干扰观察

实验值只是为了看清过程。不要把410s、并行度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]

这里故意让10011002各出现两次。

这样可以同时观察两个数字:

输入记录数: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=PAID

Flink 内部的RowData还带着动作类型:

+I INSERT -U UPDATE_BEFORE +U UPDATE_AFTER -D DELETE

JDBC 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 时发生了什么

逐条展开:

到达顺序输入batchCountreduceBuffer
1+I[1001,80]11001 -> UPSERT 80
2+U[1001,100]21001 -> UPSERT 100
3+I[1002,50]31001 -> UPSERT 1001002 -> UPSERT 50
4-D[1002]41001 -> UPSERT 1001002 -> 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()启动的定时线程低流量时控制可见延迟
CheckpointsnapshotState()保存状态前清空内存 Buffer
Closeclose()有界任务或正常结束时处理尾批

所以真实批次大小通常小于等于max-rows

例如配置 100 条,但只来了 7 条:

1 秒时间间隔先到 -> 7 条也会 Flush

又例如来了 40 条时 Checkpoint 到达:

Checkpoint 先到 -> 40 条也会 Flush

max-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 一次收到一个 SplitSplit 在 Source 内部逐条读成 RowData
invoke()逐条调用,所以 MySQL 一定逐条访问多次 invoke 之间会积累 Buffer
max-rows=100就一定有 100 个不同 Key统计的是输入条数,同 Key 会归并
一次 Flush 就只有一条 SQLUPSERT 和 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 Snapshot
  • ContinuousFileSplitEnumerator.java:生成并分配 FileStoreSourceSplit
  • FileStoreSourceReader.java:消费 Split,逐条向下游发出 RowData
  • GenericJdbcSinkFunction.java:每条记录调用writeRecord(),Checkpoint 时调用flush()
  • JdbcOutputFormat.java:维护batchCount、定时任务、行数 Flush、重试和 Close Flush
  • TableBufferReducedStatementExecutor.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 行为。

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

STM32F103车牌识别实战:硬件协同+规则引擎+微型CNN三阶架构

简介&#xff1a;本资源是面向嵌入式开发初学者与进阶工程师的STM32车牌识别实战项目&#xff0c;聚焦在资源受限MCU上实现端侧图像采集、处理与识别的完整闭环&#xff0c;解决边缘设备中轻量化视觉算法部署的核心难题。压缩包共236个文件&#xff0c;含38个C源文件&#xff0…

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

STM32F407寄存器编程驱动TFT LCD触摸屏:从FSMC配置到坐标校准实战

简介&#xff1a;本资源是面向嵌入式初学者与STM32进阶开发者的寄存器级实战例程&#xff0c;聚焦7英寸TFTLCD电容触摸屏模块在STM32F407平台上的底层驱动与交互验证。通过纯寄存器编程&#xff08;非HAL/标准外设库&#xff09;&#xff0c;完整实现LCD显示控制、CTP坐标采集、…

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

C#原生USB HID通信骨架:Windows API手撸DeviceIoControl实现

简介&#xff1a;这是一份面向C#初学者与嵌入式USB通信开发者的USB HID上位机实战源码&#xff0c;聚焦HID设备的数据收发、设备枚举与报告解析等核心问题&#xff0c;适用于工业控制、智能硬件调试、自定义HID外设联调等场景。压缩包共98个文件&#xff0c;含32个C#源码文件&a…

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

MFC Tab控件开发:从CTabCtrl到可维护Tab组件的工程实践

简介&#xff1a;本资源是一份面向MFC初学者与中级开发者的Tab Control定制化实现源码包&#xff0c;聚焦于多页界面开发中的核心控件封装与扩展实践。资源提供完整的Tabsheet类实现&#xff0c;通过继承CWnd并封装CTabCtrl&#xff0c;解决了标准MFC选项卡控件缺乏视图管理、样…

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

Arduino控制SG90 360度连续旋转舵机:脉宽、校准与开环控制

拿到一个丝印上写着“SG90”的舵机&#xff0c;如果它是 360 度连续旋转版本&#xff0c;请立刻忘掉“标准舵机是按角度转动”的习惯。你会发现servo.write(90)并不能让它停在中间&#xff0c;servo.write(0)也不会让它转 180 度后再停下来。它可能一直转&#xff0c;也可能在某…

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

小迪安全学习笔记-Day3 拓展模式以及测试会遇到的问题

Day3 拓展模式以及测试会遇到的问题WAF&#xff1a;Web应用防火墙&#xff0c;会对出入站流量进行过滤&#xff0c;安全测试手法遭到拦截。绕防火墙一般是绕很拉的防火墙&#xff0c;或者本身技术很高超能绕过。正常情况很难绕过的&#xff0c;从其他方面入手。CDN&#xff1a;…

作者头像 李华