前四篇已经把 Paimon 出仓到 MySQL 的链路拆开:
从零跑通
逐条进入与批量 Flush
UPSERT / DELETE 与 Flush 边界
Checkpoint 重放与幂等恢复
最后一篇回到生产现场。
很多任务在测试环境只有每秒几十条,看起来一切正常;上线以后更新比例、热点 Key、Checkpoint、MySQL 锁和网络延迟一起出现,才发现:
Batch 调大了,Checkpoint 变慢
并行度调高了,MySQL 连接和锁等待暴涨
Flush 间隔调长了,业务认为数据延迟
MySQL 变慢了,Flink 整条链路开始背压
无状态重跑后,目标表还残留旧数据
这一篇不提供一个“万能最佳配置”,而是给出一套能测量、能解释、能回滚的调优方法。
一、生产调优先回答三个问题
不要一上来就问:
max-rows 配多少最好?
先明确目标:
| 允许的端到端延迟是多少 | P95 小于 3 秒 |
| 峰值变化量是多少 | 5 万条/秒,更新占 70% |
| MySQL 能承受多少并发写 | 20 条连接,目标 QPS 2 万 |
没有 SLA、峰值和数据库容量,任何 Batch Size 都只是猜测。
二、先认识五个相互牵制的旋钮
这条链路最常调的不是一个参数,而是五个:
| sink.buffer-flush.max-rows | 单批输入上限、吞吐、重试成本 |
| sink.buffer-flush.interval | 低流量可见延迟、定时 Flush 频率 |
| sink.parallelism | JDBC 连接数、并发 Batch、Key 分布 |
| sink.max-retries | 短暂故障容忍与重复执行窗口 |
| Checkpoint Interval | 恢复窗口、强制 Flush 频率、Checkpoint 成本 |

这些参数相互影响。
例如缩短 Checkpoint 间隔,会更频繁地强制 Flush;即使 max-rows 很大,真实 Batch 也可能长期偏小。
三、Batch Size 调大,收益和成本分别是什么
可能收益
- 减少 executeBatch() 次数;
- 减少网络往返;
- 更充分利用 Driver Batch Rewrite;
- 降低单条语句固定开销;
- 提高高吞吐场景的总体写入效率。
可能成本
- 单批在内存中停留更久;
- 单次执行耗时更长;
- 锁持有和日志写入更集中;
- 失败后重试的数据更多;
- Checkpoint 到来时可能被一个大 Flush 拖慢;
- 大字段时内存和网络包明显增加。
所以 Batch Size 不是越大越好,而是找到收益开始变缓、风险尚可接受的区间。
四、max-rows 统计输入,不是最终 DML
回顾上一篇:
batchCount 统计进入 Sink 的记录数
reduceBuffer 按 Key 只保留最后动作
同样 1000 条输入:
| 1000 个不同 Key | 接近 1000 | 较高 |
| 10 个热点 Key 各更新 100 次 | 接近 10 | 较低,但热点竞争可能高 |
| 更新和删除混合 | 取决于每个 Key 最后动作 | 分成 UPSERT / DELETE 两组 |
因此压测报告必须同时记录:
输入条数
不同 Key 数
最终 UPSERT 数
最终 DELETE 数
只写“每秒 10 万条”无法推导 MySQL 实际负载。
五、Flush Interval 决定低流量延迟下限
低流量任务可能很久都攒不满 max-rows。
这时主要由:
'sink.buffer-flush.interval' = '1s'
控制可见延迟。
如果设置 5 秒,一条刚错过上次 Flush 的记录,可能在 Buffer 中等待接近 5 秒,除非行数或 Checkpoint 先触发。
Interval 太小
大量很小的 Batch
网络往返增多
MySQL 每秒执行次数上升
Driver Rewrite 收益下降
Interval 太大
低流量数据可见延迟升高
Buffer 驻留更久
故障时待处理范围更大
业务误判同步卡住
调整 Interval 时要看业务的延迟 SLA,不只是 Sink 吞吐。
六、并行度不是越大越快
sink.parallelism=8 通常意味着最多有 8 个 Sink 子任务并发工作,每个子任务有独立 Buffer 和 JDBC 连接。
并行度增加可能带来:
更多并发 executeBatch
更高潜在吞吐
更均匀地消化上游流量
也可能带来:
更多 MySQL 连接
更多并发事务和锁竞争
更多 InnoDB Buffer Pool 抖动
更多 Binlog / Redo 写入压力
Checkpoint 时多个子任务同时 Flush
如果 MySQL 只允许 30 条业务连接,而出仓作业开 24 条,再叠加应用连接池,很容易把连接资源耗尽。
七、先算清楚连接预算
可以先做一个简单预算:
MySQL max_connections
– 业务服务保留连接
– 运维和监控保留连接
– 其他 ETL / 查询连接
= 出仓作业可用连接预算
再分配给多个作业:
作业 A sink.parallelism
+ 作业 B sink.parallelism
+ 作业 C sink.parallelism
<= JDBC 写入连接预算
不要只看一个作业的 Web UI 决定并行度。
八、Key 倾斜会让子任务冷热不均
即使并行度为 8,也可能出现:
子任务 0:每秒 2 万条
子任务 1:每秒 500 条
其他子任务:每秒 1000 条
热点 Key 或上游分区方式会让少数子任务快速攒满并频繁 Flush,其他子任务主要靠时间 Flush。
表现为:
- 某个 Sink 子任务持续 Busy;
- 单个子任务反压明显;
- MySQL 某些索引页锁竞争;
- 全局吞吐提不上去;
- 增加并行度收益很小。
调优前要先确认瓶颈是总容量不足,还是数据倾斜。
九、Checkpoint 为什么可能制造写入尖峰
Checkpoint 到来时,每个 Sink 子任务都会执行:
outputFormat.flush();
假设 8 个子任务都积累了一批数据:
Checkpoint Barrier 到达
8 个子任务相近时间 Flush
MySQL 瞬时出现 8 个并发 Batch
如果平时主要靠较长 Interval 攒批,Checkpoint 可能形成周期性流量尖峰。
常见现象:
每到 Checkpoint 时间,MySQL QPS 和延迟出现锯齿
Sink snapshot duration 突然升高
Checkpoint 超时后作业重启,进一步放大压力
十、Checkpoint Interval 该怎样权衡
间隔更短
优点:
- 故障后重放窗口更小;
- Source 进度更频繁持久化;
- 内存 Buffer 更频繁清空。
成本:
- Checkpoint 开销更频繁;
- 更容易形成小 Batch;
- MySQL Flush 尖峰更密集;
- 状态后端和存储压力增加。
间隔更长
优点:
- Checkpoint 固定开销更低;
- Batch 更有机会积累。
成本:
- 故障重放窗口更大;
- 单次 Checkpoint 可能处理更多积压;
- Consumer 进度更新更慢;
- 数据恢复时间可能增加。
最终要结合恢复点目标和 MySQL 写入能力测试。
十一、MySQL 变慢为什么会让整条 Flink 链路反压
普通 JDBC Sink 的 Flush 是同步调用。
executeBatch() 没返回
Sink 线程无法继续处理新记录
Sink 输入队列逐渐堆积
上游算子无法继续发送
反压一路传回 Paimon Source
反压不是单独的错误,而是流系统在下游变慢时的自然流量控制。
它既保护 MySQL 不被无限请求淹没,也说明当前端到端容量不足。
十二、看到反压后不要只调大 Buffer
调大 Buffer 只能延迟问题暴露,不能让慢 MySQL 自动变快。
正确排查顺序:
1. MySQL executeBatch 延迟是否升高
2. 是否存在锁等待、死锁或慢磁盘
3. 连接数是否耗尽
4. Binlog、Redo、主从复制是否成为瓶颈
5. 是否有热点 Key / 热点索引页
6. Batch 是否过大导致单次执行过慢
7. Sink 并行度是否超出数据库承载能力
只有确认固定网络开销占主导时,继续增大 Batch 才更可能有效。
十三、MySQL 侧至少观察什么
| 连接 | 当前连接、连接失败、连接创建频率 |
| 执行 | DML QPS、Batch 延迟、影响行数 |
| 锁 | 行锁等待、死锁、长事务 |
| 存储 | Buffer Pool、磁盘 IOPS、Redo / Binlog |
| 复制 | 主从延迟、Relay Log 积压 |
| 资源 | CPU、内存、网络吞吐 |
不要只观察 MySQL CPU。
CPU 不高但锁等待严重时,继续增加并行度通常只会让情况更差。
十四、Driver Batch Rewrite 要核对什么
Flink 3.3.0 MySQL Dialect 在 URL 未显式配置时会补:
rewriteBatchedStatements=true
它允许 Connector/J 在适用时重写批处理,减少往返。
上线前仍要核对:
- 实际生效的 JDBC URL;
- Driver 版本;
- SQL 类型是否满足重写条件;
- 慢日志和网络请求形态;
- Rewrite 前后吞吐与延迟。
不要只因为配置字符串存在,就假定性能一定提高相同倍数。
十五、怎样做一轮有结论的调参实验
一次只改变一个主变量。
例如固定:
数据集
不同 Key 比例
更新 / 删除比例
Sink 并行度
Checkpoint Interval
MySQL 配置
只调整 Batch Size:
| A | 100 | 待测 | 待测 | 待测 | 待测 | 待测 |
| B | 500 | 待测 | 待测 | 待测 | 待测 | 待测 |
| C | 1000 | 待测 | 待测 | 待测 | 待测 | 待测 |
| D | 5000 | 待测 | 待测 | 待测 | 待测 | 待测 |
当吞吐增长已经很小,Checkpoint 和 MySQL P95 却明显恶化时,就不应该继续盲目调大。
十六、推荐的渐进式调优顺序
第一步:并行度 1,确认语义和正确性
第二步:固定并行度,寻找合适 Batch Size
第三步:固定 Batch,寻找满足 SLA 的 Flush Interval
第四步:逐步增加并行度,观察连接、锁和吞吐收益
第五步:调整 Checkpoint Interval,观察恢复窗口和写入尖峰
第六步:做故障恢复和重放验证
第七步:在峰值数据分布下持续压测
每一步都保留上一组稳定配置,便于回滚。
十七、端到端观测要把三层串起来

Paimon / Source 层
观察:
是否持续发现新 Snapshot
Split 是否正常分配和消费
Source 是否因为下游反压变慢
Consumer 进度是否持续推进
Flink / JDBC Sink 层
观察:
输入和输出速率
Busy / Backpressure
Checkpoint Duration 与失败原因
Task 重启次数
JDBC 连接和 Flush 异常
MySQL 层
观察:
连接、DML、锁、存储、复制和资源
只有三层时间线对齐,才能判断根因在哪一层。
十八、行数相同为什么仍然可能数据不一致
源端和目标端都是 100 万行,不代表完全一致。
可能出现:
Paimon 少 1001,多 2001
MySQL 多 1001,少 2001
总行数仍然相同
或者主键相同,但金额和状态不同。
因此对账至少分三层:
| 总量 | COUNT、分区或时间窗口行数 |
| Key 集合 | Paimon 有而 MySQL 无;MySQL 有而 Paimon 无 |
| 字段内容 | 相同 Key 的关键字段和 Hash 是否一致 |
十九、无状态重跑后为什么必须查孤儿 Key
latest-full 重新输出的是 Paimon 当前仍存在的行。
它可以:
补齐 MySQL 缺失 Key
覆盖 MySQL 中已有 Key 的旧字段值
但不能自动:
删除 MySQL 中存在、Paimon 当前不存在的孤儿 Key
重建链路时更稳妥的方案:
- 写入一张空目标表;
- 写入带版本后缀的影子表,校验后切换;
- 对账 Key 集合后再清理孤儿行;
- 尽量从有效 Checkpoint / Savepoint 恢复;
- 保持 Consumer ID 与 Snapshot 保留策略稳定。
二十、Schema 和数据类型也属于生产稳定性
上线前逐字段核对:
Paimon 类型
Flink JDBC 逻辑类型
MySQL 物理类型
重点检查:
- DECIMAL 精度和小数位;
- VARCHAR 长度和字符集;
- TIMESTAMP 精度和时区;
- NULL / NOT NULL;
- BOOLEAN 映射;
- 主键字段类型;
- 新增列、删除列和默认值。
Schema 变更不能只改 Paimon 或 Flink DDL,还要评估 MySQL 物理表和正在运行的作业是否兼容。
二十一、生产变更怎样降低风险
不要一次同时改四个参数
如果同时改变 Batch、Interval、并行度和 Checkpoint,出了问题很难定位原因。
先灰度或影子写
让少量表、少量流量或独立目标表先运行,验证语义和容量。
保留回滚点
调整前记录:
原配置
Savepoint / Checkpoint 状态
MySQL 表结构
关键监控基线
明确回滚动作
例如:
并行度过高 -> 恢复原并行度
Batch 过大 -> 恢复原 max-rows
Schema 不兼容 -> 切回旧目标表或旧作业
数据错误 -> 停止 Writer,先冻结现场和对账
二十二、一个简单的容量估算例子
假设:
峰值输入:20,000 条/秒
同 Key 归并后:约 12,000 个动作/秒
MySQL 单连接稳定处理:2,000 个动作/秒
目标利用率:不超过 70%
粗略需要的连接数:
12,000 / 2,000 / 0.7 ≈ 8.6
可以从并行度 8 或 10 开始压测,而不是直接设置 50。
但这只是起点。真实能力还受:
行大小
索引数量
更新比例
锁竞争
磁盘和 Binlog
复制延迟
影响,必须通过真实数据分布验证。
二十三、常见现象与调优方向
| MySQL QPS 很高但吞吐低 | Batch 太小、Interval 太短 | 增大 Batch 或 Interval,核对 Rewrite |
| P95 延迟过高但吞吐够 | Interval 太长、Batch 过大 | 缩短 Interval,检查 Checkpoint |
| Checkpoint 周期性变慢 | Barrier 到达时集中 Flush | 看 Sink snapshot 与 MySQL Batch 延迟 |
| 并行度升高后吞吐不升 | MySQL 已饱和或锁竞争 | 降低并发,查锁和磁盘 |
| 单个子任务持续背压 | Key 倾斜 | 查分区、热点 Key 和单 Key 写入 |
| 故障后审计重复 | 副作用非幂等 | 增加去重或改审计设计 |
| 无状态重跑仍有旧行 | 目标孤儿 Key | 做 Key 差集和清理 |
| 主从延迟升高 | Binlog / Replica 跟不上 | 降低写入峰值,评估复制容量 |
二十四、上线前最终检查清单
正确性
[ ] 两边主键与 Paimon Key 语义完全一致
[ ] 新增、更新、删除均已验证
[ ] Trigger 和副作用幂等
[ ] 无状态重建有孤儿 Key 方案
性能
[ ] 用真实 Key 分布和更新比例压测
[ ] Batch Size 已找到收益拐点
[ ] Flush Interval 满足延迟 SLA
[ ] 并行度未超过连接和锁预算
[ ] Checkpoint 不会周期性压垮 MySQL
稳定性
[ ] Checkpoint 持续成功
[ ] 故障重放实验通过
[ ] MySQL 限流、超时和重试行为明确
[ ] 监控覆盖 Source、Sink、MySQL 三层
[ ] 有可执行的回滚步骤
数据治理
[ ] 总量、Key 集合、关键字段都能对账
[ ] Schema 变更流程明确
[ ] Consumer ID 与 Snapshot 保留策略稳定
[ ] 密码和连接信息未硬编码进公开文件
二十五、整个出仓到 MySQL 子系列的最终结论
把第 4~8 篇压缩成一条完整链路:
Paimon 通过 Snapshot 和 Split 持续提供变化
↓
RowData 一条条进入 Flink JDBC Sink
↓
每个 Sink 子任务独立按主键归并并攒批
↓
行数、时间、Checkpoint 或 Close 触发 Flush
↓
JDBC 使用 UPSERT Batch 和 DELETE Batch 写 MySQL
↓
普通 SQL JDBC Sink 允许失败重放
↓
真实主键、幂等 DML 和无副作用设计让最终结果收敛
↓
Batch、Interval、并行度和 Checkpoint 共同决定生产性能
↓
三层监控与 Key 对账保证长期可运营
最后真正需要记住的不是某一组默认参数,而是四个边界:
逐条进入不等于逐条访问 MySQL
批量执行不等于与 Checkpoint 原子提交
主表结果幂等不等于所有副作用幂等
全量 UPSERT 不等于自动清理目标孤儿行
只要这四条边界清楚,遇到吞吐、延迟、重复、删除、恢复和对账问题时,就能沿着正确层次定位,而不是只在某一个参数上反复试错。
本篇关键配置与资料位置
- JdbcConnectorOptions.java:max-rows、Flush Interval、Retries 默认值
- JdbcOutputFormat.java:行数、时间、Checkpoint、Close Flush 与同步执行
- TableBufferReducedStatementExecutor.java:输入条数与最终 Key 动作数的差别
- MySqlDialect.java:MySQL Batch Rewrite 默认 URL 属性
- GenericJdbcSinkFunction.java:Checkpoint 时强制 Flush
- ContinuousFileSplitEnumerator.java:Paimon Source Snapshot 与 Split 进度
- Flink Web UI 与 MySQL Performance Schema:端到端定位的重要观测入口
本文基于 Apache Paimon 1.4.2、Apache Flink 1.20.1、Flink JDBC Connector 3.3.0-1.20 和 MySQL 8.x。文中的容量数字仅用于演示估算方法,不能替代真实环境压测。



