flink CDC任务监听mysql数据。只要不是从savepoint/checkpoint中恢复,都是能成功运行并监听数据的,但是只要从savepoint/checkpoint中恢复作业,就会报如下错误:
Caused by: java.lang.RuntimeException: One or more fetchers have encountered exception
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcherManager.checkErrors(SplitFetcherManager.java:261)
at org.apache.flink.connector.base.source.reader.SourceReaderBase.getNextFetch(SourceReaderBase.java:169)
at org.apache.flink.connector.base.source.reader.SourceReaderBase.pollNext(SourceReaderBase.java:131)
at org.apache.flink.streaming.api.operators.SourceOperator.emitNext(SourceOperator.java:419)
at org.apache.flink.streaming.runtime.io.StreamTaskSourceInput.emitNext(StreamTaskSourceInput.java:68)
at org.apache.flink.streaming.runtime.io.StreamOneInputProcessor.processInput(StreamOneInputProcessor.java:65)
at org.apache.flink.streaming.runtime.tasks.StreamTask.processInput(StreamTask.java:550)
at org.apache.flink.streaming.runtime.tasks.mailbox.MailboxProcessor.runMailboxLoop(MailboxProcessor.java:231)
at org.apache.flink.streaming.runtime.tasks.StreamTask.runMailboxLoop(StreamTask.java:839)
at org.apache.flink.streaming.runtime.tasks.StreamTask.invoke(StreamTask.java:788)
at org.apache.flink.runtime.taskmanager.Task.runWithSystemExitMonitoring(Task.java:952)
at org.apache.flink.runtime.taskmanager.Task.restoreAndInvoke(Task.java:931)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:745)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:562)
at java.lang.Thread.run(Thread.java:750)
Caused by: java.lang.RuntimeException: SplitFetcher thread 0 received unexpected exception while polling the records
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.runOnce(SplitFetcher.java:165)
at org.apache.flink.connector.base.source.reader.fetcher.SplitFetcher.run(SplitFetcher.java:114)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
... 1 more
Caused by: org.apache.kafka.connect.errors.ConnectException: An exception occurred in the change event producer. This connector will be stopped.
at io.debezium.pipeline.ErrorHandler.setProducerThrowable(ErrorHandler.java:50)
at com.ververica.cdc.connectors.mysql.debezium.task.context.MySqlErrorHandler.setProducerThrowable(MySqlErrorHandler.java:85)
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource$ReaderThreadLifecycleListener.onCommunicationFailure(MySqlStreamingChangeEventSource.java:1544)
at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:1079)
at com.github.shyiko.mysql.binlog.BinaryLogClient.connect(BinaryLogClient.java:631)
at com.github.shyiko.mysql.binlog.BinaryLogClient$7.run(BinaryLogClient.java:932)
... 1 more
Caused by: io.debezium.DebeziumException: Client requested master to start replication from position > file size Error code: 1236; SQLSTATE: HY000.
at io.debezium.connector.mysql.MySqlStreamingChangeEventSource.wrap(MySqlStreamingChangeEventSource.java:1488)
... 5 more
Caused by: com.github.shyiko.mysql.binlog.network.ServerException: Client requested master to start replication from position > file size
at com.github.shyiko.mysql.binlog.BinaryLogClient.listenForEventPackets(BinaryLogClient.java:1043)
版权声明:本文内容由阿里云实名注册用户自发贡献,版权归原作者所有,阿里云开发者社区不拥有其著作权,亦不承担相应法律责任。具体规则请查看《阿里云开发者社区用户服务协议》和《阿里云开发者社区知识产权保护指引》。如果您发现本社区中有涉嫌抄袭的内容,填写侵权投诉表单进行举报,一经查实,本社区将立刻删除涉嫌侵权内容。
这个报错的根因很明确:Error 1236: Client requested master to start replication from position > file size —— checkpoint/savepoint 里记录的 binlog 位点(文件名+offset)在恢复时已经不存在了。Flink CDC(Debezium)恢复作业时会从记录的 binlog 位点继续读,但那个 binlog 文件已经被 MySQL 清理掉(或被 reset),请求的复制起点超过了现存文件的大小,连接被 MySQL 直接拒绝。
先排查确认(在源库执行):
SHOW BINARY LOGS;
SHOW VARIABLES LIKE 'binlog_expire_logs_seconds';
对照 checkpoint 里记录的 binlog 文件名(checkpoint meta 里可以搜到),大概率已经不在 SHOW BINARY LOGS 的列表里了。
常见触发场景:作业停了几天再恢复,期间 binlog 到期被自动清理(binlog_expire_logs_seconds 设得太短)、磁盘满被手工清理、主从切换时执行过 RESET MASTER。
解决办法:
补充一句:这类位点失效本质上是自建 MySQL binlog 运维的锅——清理策略、备份、磁盘水位都要自己盯。如果不想继续跟 binlog 斗智斗勇,可以用托管版 MySQL(比如阿里云云数据库 RDS),binlog 清理和备份策略平台托管,Flink CDC 连 RDS 也是标准姿势。继续用自建也完全没问题,把保留时长调大即可。
Flink CDC 任务从 Savepoint 或 Checkpoint 恢复失败,通常涉及状态不匹配、资源瓶颈、连接器限制或配置错误。
详情参考:https://www.aliyun.com/solution/tech-solution/flink-cdc-realize-data-synchronization?source=5176.29345612&userCode=e6tbwq9f
为了帮你精准定位问题,我将常见的错误场景及排查步骤梳理如下:
核心排查:状态与元数据匹配问题
这是最常见的原因。Flink 依赖算子的 UID 来映射状态。
错误现象:报错 Cannot map checkpoint/savepoint state for operator xxx to the new program 或 Failed to rollback to checkpoint/savepoint。
根本原因:
代码逻辑变更:如果你修改了作业代码(如增删算子、改变拓扑结构),且没有为算子指定固定的 UID,Flink 会自动生成新的 UID,导致无法从旧状态中恢复。
连接器版本变更:例如从旧版 FlinkKafkaConsumer 迁移到新版 KafkaSource,或者升级了 Flink CDC 连接器版本,两者的状态数据结构不兼容。
解决方案:
固定 UID:在代码中显式调用 .uid("unique_id"),确保重启前后 ID 一致。
检查版本:确保恢复作业时使用的 Flink 引擎版本和连接器版本与创建 Savepoint 时一致。
资源瓶颈:内存溢出 (OOM) 与 RPC 超限
当 Savepoint 的元数据文件(_metadata)过大时,恢复过程会消耗大量资源。
错误现象:
java.lang.OutOfMemoryError: Java heap space(JobManager 内存溢出)。
The rpc invocation size ... exceeds the maximum akka framesize(RPC 调用超限)。
根本原因:
Savepoint 中包含大量的小状态(如 Kafka 分区偏移量),导致 _metadata 文件膨胀(从几 KB 涨到几百 MB)。
JobManager 在加载该文件时堆内存不足,或通过 RPC 分发状态时超过了 Akka 默认的 64MB 限制。
解决方案:
调大内存:增加 JobManager 的堆内存配置(如 jobmanager.memory.process.size)。
调大 RPC 限制:在 flink-conf.yaml 中增加 akka.framesize 的值(例如 104857600b 即 100MB)。
Flink CDC 特有的恢复限制
CDC 连接器在全量与增量阶段的切换点比较特殊,容易踩坑。
错误现象:恢复后无数据输出、作业频繁重启或报错。
常见场景与对策:
全量阶段修改表结构:如果在 Savepoint 后,对源表进行了 DDL 操作(如加字段),或者使用了 pt-osc 等工具变更表结构,恢复后可能无法解析。
对策:升级 Flink CDC 版本(如 VVR 11.2+),并开启 scan.parse.online.schema.changes.enabled: true。
全量阶段增删表:在全量快照阶段保存 Savepoint,然后修改了 tables 配置(增加或减少了同步的表),再恢复作业。
对策:不支持在全量阶段修改表列表后恢复。必须等全量阶段结束进入增量阶段后,再修改配置并开启 scan.newly-added-table.enabled: true。
Binlog 位点过期:如果 Savepoint 保存时间太久,对应的 Binlog 文件在 MySQL 端已被清理。
对策:检查 MySQL 的 expire_logs_days 配置,确保 Binlog 保留时间足够长。
具体的排查步骤建议
建议你按照以下顺序进行“体检”:
查看 JobManager 日志:
搜索关键词 Exception、Error、Caused by。
如果是 OutOfMemoryError,直接调大 JobManager 内存。
如果是 Cannot map...,检查代码中的算子 UID 是否变动。
检查 Savepoint 目录:
查看 Savepoint 目录下的 _metadata 文件大小。如果超过几十 MB,说明状态过于臃肿,建议优化状态后端或清理无用状态。
验证 CDC 配置:
确认 scan.startup.mode 是否正确。如果是从 Savepoint 恢复,通常不需要指定该参数,Flink 会自动读取状态中的位点。
确认 server-id 是否冲突。如果多个作业使用相同的 server-id 连接同一个 MySQL,会导致连接被踢下线。
尝试“无状态”重启(最后手段):
如果状态文件已损坏且无法修复,只能放弃 Savepoint,使用 --allowNonRestoredState 参数(忽略无法恢复的状态)或直接从最新位点(latest-offset)重新启动。注意:这可能会导致数据丢失或重复。
验证 Binlog 文件是否存在
在尝试恢复 Flink 任务前,你可以先去 MySQL 服务器上确认 Flink 需要的 binlog 文件是否还存在。
查看 Flink 需要的 binlog 位置: 这个信息通常在 Flink 的错误日志中可以找到,或者在 Flink UI 的 Checkpoint 详情里。
查看 MySQL 服务器上可用的 binlog 文件: 在 MySQL 中执行:
SQL
SHOW BINARY LOGS;
这个命令会列出所有当前可用的 binlog 文件。检查一下 Flink 需要的文件是否在这个列表里。如果不在,就说明它已经被清理了。
这是一种下策,只有在你可以接受数据丢失(从任务停止到重启这段时间的数据)的情况下才能使用。
当你确定无法从旧的 savepoint 恢复时,你可以选择放弃这个 savepoint,然后重新启动一个新任务,并设置启动参数从最新的位置开始消费。
Java
// Flink SQL
'scan.startup.mode' = 'latest-offset'
// DataStream API
MySqlSource.builder()
.startupOptions(StartupOptions.latest())
// ... 其他配置
.build();
⚠️ 警告:这样做会导致停机期间的所有数据变更全部丢失,请务必谨慎操作!
实时计算Flink版是阿里云提供的全托管Serverless Flink云服务,基于 Apache Flink 构建的企业级、高性能实时大数据处理系统。提供全托管版 Flink 集群和引擎,提高作业开发运维效率。