ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Flink自定义MySQL Source开发实战:RichSourceFunction与状态恢复

2026/9/9 21:46:19 拓冰建站 浏览量
Flink自定义MySQL Source开发实战:RichSourceFunction与状态恢复 做 Flink 开发的朋友迟早会遇到这么一个问题明明内置了 JDBC connector也能读 MySQL为什么我还要自己写一个 Source坦白说大部分场景下直接用 Flink JDBC connector 或 Flink CDC 就够了。但当你遇到分库分表读取、按业务字段轮询、字段裁剪、多表关联后做广播流、或者需要对接老旧的 JDBC 驱动比如某些国产数据库的方言时你大概率就得自己动手写自定义 Source 了。这篇就完整拆一遍 Flink MySQL Source 自定义开发的全过程包括接口选型、核心代码实现、并行度设计、状态恢复、以及我实际踩过的坑。1. 整体设计思路为什么自己写以及怎么写1.1 先搞清楚内置连接器和自定义 Source 的边界很多人一上来就写代码其实先做方案对比更省时间。Flink 官方提供了两种读取 MySQL 的现成方案Flink JDBC Connectorflink-connector-jdbc适合做维表关联、批量读取、写入通过 JDBC 驱动以批的方式拉取数据。Flink CDC Connectors如 flink-cdc-connector基于 binlog 实现增量同步适合实时性要求很高的场景。既然有现成的为什么还要自定义我把实际工作中会遇到的情况列一下场景内置 JDBC ConnectorFlink CDC自定义 Source分钟级准实时轮询可以但不太好控制大材小用且 binlog 权限难申请最合适轮询间隔可控分库分表合并读取配置麻烦可以改表名单但多表匹配逻辑弱代码里随意拼表名、加判断需要做字段解密、格式转换后输出要用 SQL 包一层也要包一层Source 内直接处理依赖老驱动/老版本 MySQL兼容性看版本要求 MySQL 5.7 且开启 binlog任何 JDBC 驱动都行学习 Flink Source 原理黑盒黑盒提升明显自定义 Source 不是用来替代 CDC 的它的定位更多是“轻量、可控、满足特定业务逻辑”。我这次的需求其实很典型上游有一个老系统MySQL 库里有一张流水表需要每分钟把新增的数据同步到数仓做轻度汇总。不能申请 binlog 权限也不能接受太重的方式于是自定义一个轮询读取的 MySQL Source 就成了最优解。1.2 自定义 Source 的三种实现方式选型Flink 开发中实现 Source 有几种方式选错了后面会很痛苦SourceFunction最简单但无法感知并行度、无状态管理、checkpoint 时不会保存读取进度。适合演示不适合生产。RichSourceFunction继承了 AbstractRichFunction可以拿到 RuntimeContext能访问状态、广播流能做 checkponit。SourceReaderFLIP-27 新接口基于异步拉取模型配合 SplitReader适合实现 Kafka 那种高吞吐的场景。但代码复杂度高一个最小实现都要好几个类。对于 MySQL 轮询这种场景RichSourceFunction 是最优选择。理由很简单我们需要在并行子任务中各自读取不同分片的数据需要记录“读到哪一行了”并把这个进度保存到 Flink 的状态后端里。RichSourceFunction 支持 operator stateSourceReader 虽然更强但重写成本太高没必要。1.3 核心执行模型Source 的运行循环自定义 MySQL Source 的运行逻辑其实很朴素就是一个循环初始化open 创建数据库连接池 读取恢复状态如果从 checkpoint 恢复 确定当前子任务的分片范围 运行时run 循环执行 按分片范围和上次位置查询数据 将每条数据发送到下游 记录最新位置到状态 根据业务情况 sleep 一段时间轮询间隔 检查取消标志如果收到 cancel 则退出 清理cancel/close 关闭查询资源 关闭连接池这个模型和读取消息队列是类似的都是“持续产生数据流”。只不过 MySQL 没有“offset”概念我们要用业务字段一般用自增主键或时间戳来模拟 offset实现断点续读。这一块的细节会在后面代码里展开。2. 核心细节解析与实操要点2.1 数据读取策略增量轮询 vs 全量扫描自定义 MySQL Source 最核心的问题是“怎么查”。我见过不少新手上来就把整张表 select * 出来然后每轮都全量发往下游。这在小数据量时没问题但数据量到百万级就会出大问题——每轮全量扫描把 MySQL 的 IO 打满下游还收到大量重复数据。我的做法是两种模式结合增量模式表里有自增主键或时间戳字段时记录上次的最大值本次查询带 where 条件只读新增的数据。全量模式首次启动时没有断点或者业务需要全量刷新时一次性全表扫描。增量模式的 SQL 大致长这样SELECT * FROM biz_order WHERE id ? AND id ? ORDER BY id ASC后面两个参数分别对应“上次位置”和“当前批次上限”。这里的 id 是主键用主键做断点最可靠因为索引走的就是主键查询效率最高。如果主键不是自增的而是雪花 ID也可以用但 where 条件变成id ?后要保证排序稳定否则会出现漏数据。如果表没有主键只有创建时间字段那就用时间戳SELECT * FROM biz_order WHERE create_time ? AND create_time ? ORDER BY create_time ASC但用时间字段有个坑同一秒内可能有多条数据且两个查询之间数据持续写入。如果每轮读一次上一轮读到2025-01-01 10:00:00下一轮查 10:00:00中间那条 10:00:00.500 的记录就漏掉了。解决方法是把轮询窗口做成“重叠”的比如每次从上次时间 - 5秒开始查然后在上游数据表结构允许的情况下加一个op_time处理标记否则就要在数据落库时设计一个数据批次号batch_id字段比如每 5 分钟一个批次这样既好查也好去重。我在实际项目中更推荐加 batch_id 的方案因为时间戳方案在处理“跨秒边界”的问题上永远会多出边际成本。2.2 状态存储与断点续读原理自定义 Source 想做到“任务重启不丢数据”就必须配合 Flink 的 checkpoint 机制。RichSourceFunction 里有一个CheckpointedFunction接口实现它以后Flink 在制作 checkpoint 时会把offset保存到状态后端比如 RocksDB。具体做法public class MysqlRichSourceFunction extends RichSourceFunctionRowData implements CheckpointedFunction { private transient ListStateLong offsetState; private long offset 0L; Override public void snapshotState(FunctionSnapshotContext context) throws Exception { offsetState.clear(); offsetState.add(offset); } Override public void initializeState(FunctionInitializationContext context) throws Exception { ListStateDescriptorLong descriptor new ListStateDescriptor(mysql_offset, Long.class); offsetState context.getOperatorStateStore().getListState(descriptor); if (context.isRestored()) { offset offsetState.get().iterator().next(); } } }这里的offset在增量模式下就是“上次读到的主键最大值”在时间戳模式下就是“上次读到的最新时间”。注意一个关键点这里用的算子状态Operator State不是键控状态Keyed State。因为 Source 算子没有 keyBy无法用 keyed state。用 ListState 的好处是可以在并行度改变时均匀分配。但是要注意如果并行度从 4 改成 2原来每个任务里保存的 offset 如何合并状态恢复时的重分布逻辑是 Flink 框架自动处理的。这带来了一个隐患并行子任务各自维护 offset但 MySQL 的实际分组方式却是自定义的分片策略如果分片策略和状态恢复不对应会出现重复读或漏读。因此我更建议把“读取进度”保存为两部分当前子任务负责的分片范围 DESCRIPTOR比如task1: id 1-100000, task2: id 100001-200000当前子任务 read offset这两部分合在一起才是一个完整的“Source 位点”下次恢复时才能按原来的分片继续读。2.3 并行度和分片策略设计自定义 MySQL Source 的并行度不是随便设的。并行度过大会把 MySQL 的连接数和 IO 打爆并行度过小又起不到加速作用。这里需要用“分片”的思路来解耦并行度和 MySQL 压力如果表有自增主键可以按主键区间分片子任务 0id BETWEEN 0 AND 1000000子任务 1id BETWEEN 1000001 AND 2000000以此类推如果表没有主键就按MOD(id, parallelism)取模分片但这种方式对索引不友好数据量大时会全表扫描。我在项目中用的是“动态计算分片”的方式在 open() 方法里先查一次最大值和最小值然后按当前并行度均分Override public void open(Configuration parameters) throws Exception { int indexOfThisSubtask getRuntimeContext().getIndexOfThisSubtask(); int numberOfParallelSubtasks getRuntimeContext().getNumberOfParallelSubtasks(); // 查询表的最大、最小主键 long minId queryMinId(); long maxId queryMaxId(); long range maxId - minId 1; long step range / numberOfParallelSubtasks 1; this.startId minId step * indexOfThisSubtask; this.endId Math.min(maxId 1, startId step); this.currentId this.startId; }这个方式看起来简单但在数据分布极不均匀时会有问题主键是自增的业务高峰期写入快ID 之间可能有空洞这个无所谓因为我们是按 ID 范围读空洞不影响结果只是查询时会扫描到不存在的 IDMySQL 走主键索引扫描空范围性能可接受。真正麻烦的是主键不是连续自增还伴随表数据动态删除时按范围分片会不均衡。比如 A 段 0~1000 只剩 10 条数据B 段 5000~6000 有 900 条数据并行度虽然开了但只有 B 段的任务一直在忙。这种时候就得用“按实际行的采样值分片”的方法了比如先SELECT id FROM table ORDER BY id LIMIT 1000 OFFSET ...拿到样本点再按样本值切分区间。做起来复杂点但效果均衡。3. 实操过程与核心环节实现3.1 开发环境准备开发 Flink 自定义 Source 的环境比较基础JDK 1.8Maven 3.6Flink 1.13 或 1.14/1.15接口略有差异下面代码以 1.13 为主兼容 1.15 问题不大MySQL 驱动 8.0.x推荐用 HikariCP 连接池Maven 依赖如下如果是 1.14/1.15 把版本号换掉即可properties flink.version1.13.6/flink.version /properties dependencies dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_2.12/artifactId version${flink.version}/version scopeprovided/scope /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-core/artifactId version${flink.version}/version scopeprovided/scope /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version8.0.28/version /dependency dependency groupIdcom.zaxxer/groupId artifactIdHikariCP/artifactId version4.0.3/version /dependency /dependencies注意几点flink-streaming-java和flink-core的 scope 要设成provided避免和 Flink 集群里自带的包冲突。HikariCP 和 MySQL 驱动需要打包进 fat jar所以不去掉 scope。如果你用的是 flink-connector-jdbc 做写入就额外加一份 JDBC connector 依赖。3.2 核心代码实现完整版 RichSourceFunction下面直接给出一份生产可用的简化版代码。我这里统一以“自增主键增量轮询 分片 checkpoint 恢复”作为主逻辑大家可以根据自身情况替换查询 SQL。package com.example.flink.source; import com.zaxxer.hikari.HikariConfig; import com.zaxxer.hikari.HikariDataSource; import org.apache.flink.api.common.state.ListState; import org.apache.flink.api.common.state.ListStateDescriptor; import org.apache.flink.api.common.typeinfo.TypeInformation; import org.apache.flink.configuration.Configuration; import org.apache.flink.runtime.state.FunctionInitializationContext; import org.apache.flink.runtime.state.FunctionSnapshotContext; import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction; import org.apache.flink.streaming.api.functions.source.RichSourceFunction; import org.apache.flink.types.Row; import javax.sql.DataSource; import java.sql.Connection; import java.sql.PreparedStatement; import java.sql.ResultSet; import java.sql.SQLException; public class MysqlIncrementSource extends RichSourceFunctionRow implements CheckpointedFunction { private static final long serialVersionUID 1L; private final String jdbcUrl; private final String username; private final String password; private final String tableName; private final String pkColumn; private final long pollIntervalMs; private transient DataSource dataSource; private transient volatile boolean running true; private transient long currentId 0L; private transient long endId Long.MAX_VALUE; private transient ListStateLong offsetState; private transient ListStateLong rangeState; public MysqlIncrementSource(String jdbcUrl, String username, String password, String tableName, String pkColumn, long pollIntervalMs) { this.jdbcUrl jdbcUrl; this.username username; this.password password; this.tableName tableName; this.pkColumn pkColumn; this.pollIntervalMs pollIntervalMs; } Override public void open(Configuration parameters) throws Exception { HikariConfig config new HikariConfig(); config.setJdbcUrl(jdbcUrl); config.setUsername(username); config.setPassword(password); config.setMaximumPoolSize(5); config.setMinimumIdle(1); config.setConnectionTimeout(5000); dataSource new HikariDataSource(config); int taskIndex getRuntimeContext().getIndexOfThisSubtask(); int parallelism getRuntimeContext().getNumberOfParallelSubtasks(); // 如果状态里没有分片范围则从 MySQL 里查一次 min/max id 并计算分片 if (rangeState null || !rangeState.get().iterator().hasNext()) { long[] minMax queryMinMaxId(); long minId minMax[0]; long maxId minMax[1]; long step (maxId - minId) / parallelism 1; currentId minId step * taskIndex; endId Math.min(maxId 1, currentId step); } else { // 从状态恢复范围 java.util.ListLong range new java.util.ArrayList(); rangeState.get().forEach(range::add); currentId range.get(0); endId range.get(1); } // 恢复已经读取的位置 if (offsetState ! null offsetState.get().iterator().hasNext()) { currentId offsetState.get().iterator().next(); } // 校正起点不能超过终点 if (currentId 0) { currentId 0; } if (endId currentId) { endId currentId; } } Override public void run(SourceContextRow ctx) throws Exception { while (running currentId endId) { // 1. 查询当前批次数据 long batchStart currentId; long batchEnd Math.min(currentId 1000, endId); try (Connection conn dataSource.getConnection(); PreparedStatement ps conn.prepareStatement( SELECT * FROM tableName WHERE pkColumn ? AND pkColumn ? ORDER BY pkColumn ASC)) { ps.setLong(1, batchStart); ps.setLong(2, batchEnd); try (ResultSet rs ps.executeQuery()) { while (rs.next()) { Row row Row.withNames(); // 这里按列名取避免写死列顺序 for (int i 1; i rs.getMetaData().getColumnCount(); i) { String colName rs.getMetaData().getColumnLabel(i); row.setField(colName, rs.getObject(i)); } ctx.collect(row); } } } // 2. 更新 offset currentId batchEnd; // 3. 如果已经读到分片末尾等待下一轮轮询此时不清空当前状态以便从 maxId 继续增量 if (currentId endId) { Thread.sleep(pollIntervalMs); // 增量模式下如果读到了分片终点则把终点延伸 // 这里简化为扩展 endId 为当前 maxId下一次轮询时重新查 long[] minMax queryMinMaxId(); endId minMax[1] 1; } } } Override public void cancel() { running false; } Override public void close() throws Exception { if (dataSource ! null) { dataSource.close(); } } Override public void snapshotState(FunctionSnapshotContext context) throws Exception { offsetState.clear(); offsetState.add(currentId); rangeState.clear(); rangeState.add(currentId); rangeState.add(endId); } Override public void initializeState(FunctionInitializationContext context) throws Exception { ListStateDescriptorLong offsetDesc new ListStateDescriptor(mysql-offset, TypeInformation.of(Long.class)); ListStateDescriptorLong rangeDesc new ListStateDescriptor(mysql-range, TypeInformation.of(Long.class)); offsetState context.getOperatorStateStore().getListState(offsetDesc); rangeState context.getOperatorStateStore().getListState(rangeDesc); } private long[] queryMinMaxId() throws SQLException { try (Connection conn dataSource.getConnection(); PreparedStatement ps conn.prepareStatement(SELECT MIN( pkColumn ), MAX( pkColumn ) FROM tableName); ResultSet rs ps.executeQuery()) { rs.next(); return new long[]{rs.getLong(1), rs.getLong(2)}; } } }这段代码是能够在生产环境跑起来的最小骨架它涵盖了连接池HikariCP动态创建和关闭open 时按并行度计算分片范围checkpoint 状态保存 offset rangerun 循环按批查询、发送、更新 offset增量模式下读到末尾后扩展 endId实现持续增量读取3.3 参数计算分片步长和轮询间隔怎么定很多新手会直接抄代码但一上生产就出问题根源往往是参数拍脑袋。我建议按这几个维度来算分片步长step。分片步长不是越小越好。比如有一张 1 亿行的流水表主键 1~1 亿并行度设为 8步长为 1250 万。每次查询 1000 条的话一个批次大概扫 12500 次索引才能读完一个分片。如果读太快每轮几秒就扫完了然后要做轮询等待。这里关键要看你每批的游标查询成本WHERE id x LIMIT 1000在同一条索引上查询1000 条的耗时通常在十几毫秒内所以单任务每秒处理 1~2 万行是合理的。按这个速度一个分片 1250 万行大约需要 10 分钟。如果你的业务表每分钟新增几百行这种分片设计就绰绰有余。轮询间隔pollIntervalMs。这个参数决定了准实时性的上限。假设轮询间隔设成 1 分钟那么数据最多延迟 1 分钟才能被读到。如果要求“秒级”你就得把间隔调到 1~5 秒但这样 MySQL 会被频繁执行SELECT MAX(id)或SELECT * FROM table WHERE id ? LIMIT 1000对库的压力成倍增加。我的经验是除非业务明确要求否则轮询间隔建议 30 秒以上。因为你算一下就会发现间隔 5 秒和间隔 30 秒对用户感知的差别不大但对 MySQL 的 QPS 影响是 6 倍。每批查询行数batchSize。上面代码里写的是 1000。这个值要结合“单条 Row 的大小”来定。如果一行数据平均 500 字节1000 行也就 500 KB在网络传输和 Flink 序列化上开销很小。但如果一行数据有几 KB比如有 text 字段或 JSON 字段1000 行会产生几 MB 的数据这时候建议调整为 200~500。更好的做法是让 batchSize 由预估单行字节数反推比如目标控制每批 1 MB就用1MB / 预估行大小来定。3.4 从 MySQL 行列格式到 Flink 的数据类型转换自定义 Source 最终要把数据从 JDBC 的 ResultSet 转成 Flink 内部数据类型。上面代码用Row.withNames()是按列名映射比较灵活但要注意MySQL 的TINYINT(1)会映射成 Boolean如果你的表里有 is_deleted 这类字段没问题但如果你的字段叫 count且值是 0/1期望当整数用就得注意 Flink 类型系统会把它当成 Boolean。MySQL 的DECIMAL映射成 BigDecimalFlink 的 Row 可以直接装。MySQL 的DATETIME/TIMESTAMP映射成java.sql.Timestamp在 Flink 里处理时要转成 localdatetime 或直接用Timestamp类型注意时区问题。如果对性能要求很极致可以用RowData二进制行和GenericRowData但代码复杂度会涨不少。我建议先用 Row/RowData 打通流程后续如果要做深层优化再换。4. 常见问题与排查技巧实录4.1 连接被 MySQL 主动断开的问题这是我遇到最多的坑。MySQL 默认wait_timeout是 8 小时如果任务里跑了一个长轮询连接池里的连接空闲时间超过 wait_timeoutMySQL 会主动断开。下一次请求时虽然 HikariCP 默认做了连接测试但有些老版本的 HikariCP 不一定会检测到断开连接导致 Source 直接抛异常。解决方式分两层HikariCP 里加上connectionTestQuery或connectionInitSqlconfig.setConnectionTestQuery(SELECT 1); config.setMaxLifetime(600000); // 连接最大存活10分钟 config.setIdleTimeout(300000); // 空闲5分钟后释放 config.setValidationTimeout(3000);在 run 循环里每次查询失败时不要直接退出而是尝试重连catch (SQLException e) { LOG.warn(Query MySQL failed, retrying..., e); Thread.sleep(1000); continue; }但是注意不能无限重试掩盖真正的 SQL 错误。建议加个计数器连续失败超过 5 次就抛出异常让 Flink 任务进入失败重启流程而不是变成僵尸任务。另外Flink 任务重启之后 checkpoint 恢复时会从最近一次的 offset 继续读。这要求你的 MySQL 查询是幂等且可回放的。也就是说从同一个 offset 重复读取几次下游得到的数据应该是同一批不能因为读了两遍就把数据写重复。如果下游是 Kafka SinkKafka 默认 at-least-once 会重复如果下游是 MySQL 写入你需要在 Sink 里做 upsert。4.2 并行度调整后状态恢复错乱Flink 的 ListState 在并行度改变时会重新分配 state但分配方式是 round-robin不一定能保证原来 4 个分片的数据恢复到新的 2 个分片后还逻辑正确。原因在于ListState 恢复时是按元素粒度重新分配的不是按“分片”整体分配的。如果并行度从 4 改成 2Flink 会把 4 个任务里每个任务存的 offset 打散合并到 2 个任务里。但你的分片范围range也是存在 ListState 里的它会被拆散。于是可能发生这种情况task0 之前负责 1~10000 和 30000~40000 两个区间恢复后 task0 拿到的是 20000~30000 的区间导致漏读。我在实际项目里的做法是给每个任务存一个完整的分片状态对象比如一个 Class包含 startId 和 endId这样 ListState 里的每个元素是完整的分片而不是分散的数字。另外调整并行度时尽量遵循“状态不恢复手动指定起点”的方式尤其当你的 MySQL Source 本身就是幂等读取时稍微丢一点数据不如错乱丢一片。具体优化public static class MysqlSplit { public long startId; public long endId; public long offset; }然后ListStateMysqlSplit来保存分片。这样并行度改变时Flink 是按 split 对象分配的每个分片整体迁移逻辑才不会乱。4.3 时区问题时间字段偏移 8 小时这是 JDBC 读取 MySQL 的经典问题。MySQL 驱动的serverTimezone参数如果设置不对会导致TIMESTAMP字段返回的时间比数据库里少了 8 小时或多了 8 小时。推荐统一在 JDBC URL 里显式指定jdbc:mysql://127.0.0.1:3306/biz_db?useUnicodetruecharacterEncodingutf8useSSLfalseserverTimezoneAsia/ShanghaiuseCursorFetchtrue只要连接串统一前后端、Flink 任务里也用同一时区这个问题就能规避。另外要注意如果 Source 里读取的是字符串类型的时间通过rs.getString()拿到的就是数据库原始字符串不受时区影响但如果你用rs.getTimestamp()拿就会套 JDBC 驱动的时区逻辑。4.4 游标与流式读取大表全量扫描时的内存问题默认情况下 MySQL JDBC 驱动会把全表查出来的结果一次性加载到 JVM 内存里再交给ResultSet遍历。当表很大时JVM 直接 OOM。解决方案是在 JDBC URL 中开启useCursorFetchtrue并给 Statement 设置fetchSizeps.setFetchSize(500);useCursorFetch会让 MySQL 服务器端以游标方式逐批返回数据客户端每次只缓冲 500 行。这个选项在 MySQL 8.0 驱动上稳定但在 MariaDB 驱动上不支持需要换用setFetchSize(Integer.MIN_VALUE)这种流式读取方式。注意useCursorFetchtrue对 MySQL 5.x 的驱动也有效但必须配合setFetchSize否则默认还是全量加载。4.5 SQL 注入风险和表名/字段名拼接自定义 Source 里 SQL 往往是动态拼接的比如表名、主键名。如果这些值来自外部配置而且你直接把用户输入的字符串拼进 SQL就有注入风险。解决方式表名、字段名用白名单校验不允许动态传入。参数值一律用PreparedStatement占位符传入。不要直接在 SQL 里拼tableName除非能保证配置来源可控。我见过一个真实事故某团队把业务库名当参数传入结果业务方不小心传了带空格和分号的字符串把线上表 drop 了。虽然这是极低概率的事故但在 SI 系统里千万别开这个口子。4.6 算子状态一直很大或 checkpoint 超时如果自定义 Source 每次 checkpoint 时把大对象比如整表快照存进 statecheckpoint 就会越来越慢。对于 MySQL Sourcestate 应该只保存“位置”信息几个 long而不是数据本身。如果你发现 state 大小异常检查是不是在 snapshotState 里误存了Row对象或ResultSet。另外snapshotState是同步调用会阻塞 source 的 run 循环。如果你在 snapshotState 里访问 MySQL 查询那 checkpoint 时间会被拉得很长而且可能会因为 DB 抖动直接导致 checkpoint 失败。所以 state 必须只做内存操作不做 IO。5. 生产环境中的更多考虑5.1 下游去重与幂等设计自定义 MySQL Source 即使是增量读取也可能会因为任务重启、并行度调整、MySQL 主从切换等问题导致重复读取。所以下游最好做一层幂等处理。如果你是写入 Kafka可以在 Source 里给每行数据附加一个唯一业务键比如主键下游消费者按主键去重如果你是写入 ClickHouse、Doris 这类 OLAP 系统一般支持 upsert 语义可以接受少数重复写入。一个简单策略在 Row 里多放一个字段比如src_table_pk作为下游去重的依据。这样即使上游重复读下游也能兜底。5.2 动态感知表结构变更MySQL 表结构不是一成不变的。比如业务方在表里新增了一个字段而你的自定义 Source 建表语句或 Row 结构没有同步修改下游拿到的 Row 就会缺列或者查的时候 SQL 报错。我建议在 Source 里做一次轻量的“结构发现”private ListString queryColumns(Connection conn) throws SQLException { DatabaseMetaData metaData conn.getMetaData(); try (ResultSet rs metaData.getColumns(null, null, tableName, %)) { ListString cols new ArrayList(); while (rs.next()) { cols.add(rs.getString(COLUMN_NAME)); } return cols; } }然后运行时用列名列表动态构造SELECT col1,col2,... FROM table而不是写死SELECT *。这样新增字段时不用重启任务就能自动带出来。但要注意字段顺序变了下游解析方式也要够灵活否则按位置取值的代码会崩。5.3 延迟指标监控这个很容易被忽略。既然你是自定义 Source就要自己负责“监控”。在 run 循环里统计几个指标当前 offset / 上一次查询耗时每次查询返回的行数从上次查询到当前时间的时间差lag可以用 Flink 的RuntimeContext#getMetricGroup()注册 Gauge 或 CountergetRuntimeContext().getMetricGroup().gauge(mysql-lag-ms, () - System.currentTimeMillis() - lastQueryTime); getRuntimeContext().getMetricGroup().gauge(mysql-current-offset, () - currentId);这些指标在 Flink Web UI 的 Task Metrics 里可以直接看到也可以用 Prometheus 抓去做告警。数据延迟超过阈值了系统能第一时间通知你而不是等业务方找上门。5.4 无主键表的处理方案有些 MySQL 表没有主键也没有合适的自增列。这时做增量读取就非常麻烦。我遇到过的方案有几个用binlog的位点做增量这基本就是 Flink CDC 干的事了。在表中加一个etl_id自增列让 DBA 加一个索引由我们来维护。采用全量比对方式每轮全量读取在 Flink 端做 groupBy 去重相当于“准实时同步全量表”。方案 3 最笨但最通用适合数据量不大的表比如几十万行。它的缺点也很明显每轮全量读MySQL 压力大Flink 端要做状态管理来去重。如果真遇到无主键大表我建议和 DBA 沟通加列这是最合理的做法。6. 自定义 Source 开发的小结思考从接口选型、分片设计、状态恢复到踩坑排查自定义 Flink MySQL Source 的过程其实也是一次对 Flink 状态机制和 Source 运行模型的深度复习。很多时候我们张口闭口“状态后端”“checkpoint 语义”但到真正写一个带状态的 Source 时才会发现很多细节并不像文档里写的那么轻松。我个人在实际项目里最深的体会是宁可代码多写一点分片和恢复逻辑也不要贪图简单直接全表扫描 手动控制重启。数据量小的时候怎么折腾都不会出事数据量一旦上来没有状态恢复意识的 Source 就是一颗定时炸弹每一次重启都意味着数据缝隙和重复。后续如果要在这个基础上继续扩展可以考虑支持多表读取并合并为一个流匹配动态路由规则。增加全量增量自动切换逻辑首次启动全量之后走增量。把通用逻辑抽成一个框架通过配置化驱动不同表的读取。接入 Flink CDC 做增量阶段自定义 Source 做存量阶段二者无缝衔接。最后再分享一个实操小技巧调试自定义 Source 时不要一上来就提交到集群先用本地 IDE 跑StreamExecutionEnvironment.execute()把并行度设成 1直接看输出。这样排查 SQL 问题、类型映射问题会快十倍。等本地逻辑完全稳定了再上集群测并行度和状态恢复。