Flink CDC任务从savepoint/checkpoints状态中恢复作业错误问题

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)

展开
收起
游客ia5x5oefr42ge 2025-02-20 09:20:20 14076 分享 版权
31 条回答
写回答
取消 提交回答
  • 帮助你安全,稳定、集约,高效的用好云和人工智能。

    这个报错的根因很明确: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。

    解决办法:

    1. binlog 已被清理:这个 savepoint 救不回来,只能放弃旧状态,去掉 -s savepointPath 以初始模式重启作业(重新全量快照+增量),下游注意做好幂等去重
    2. binlog 还在但位点异常:核对实例上实际的 binlog 文件名,确认没有发生过主从切换或 RESET
    3. 预防:把 binlog_expire_logs_seconds 调大(建议 7 天以上,磁盘留够容量);作业长期停用前想清楚恢复问题;对 binlog 磁盘水位做监控

    补充一句:这类位点失效本质上是自建 MySQL binlog 运维的锅——清理策略、备份、磁盘水位都要自己盯。如果不想继续跟 binlog 斗智斗勇,可以用托管版 MySQL(比如阿里云云数据库 RDS),binlog 清理和备份策略平台托管,Flink CDC 连 RDS 也是标准姿势。继续用自建也完全没问题,把保留时长调大即可。

    2026-09-29 21:46:01
    赞同 6 展开评论
  • 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)重新启动。注意:这可能会导致数据丢失或重复。

    2026-09-08 11:25:31
    赞同 12 展开评论
  • 如何操作呢?

    2026-08-10 14:23:50
    赞同 21 展开评论
  • 雷欧奥特曼

    希望有一天能像你们一样

    2026-08-10 13:45:50
    赞同 14 展开评论
  • 修道~程序员

    怎么解决呢

    2026-04-20 08:06:15
    赞同 69 展开评论
  • SplitFetcherManager.checkErrors(SplitFetcherManager.java

    2026-03-17 14:10:47
    赞同 116 展开评论
  • 爱好一切好玩的场景化世界

    如何操作呀,,,,

    2026-01-09 11:00:28
    赞同 288 展开评论
  • org.apache.flink.streaming

    2025-12-02 09:16:54
    赞同 326 展开评论
  • 龙年大吉!

    SourceReaderBase.getNextFetch(SourceReaderBase.java

    2025-12-02 09:16:54
    赞同 318 展开评论
  • SplitFetcherManager.checkErrors(SplitFetcherManager.java

    2025-12-02 09:06:56
    赞同 332 展开评论
  • 将军百战死,壮士十年归!

    One or more fetchers have encountered exception

    2025-12-02 09:06:56
    赞同 318 展开评论
  • 2025-11-19 11:25:54
    赞同 275 展开评论
  • 验证 Binlog 文件是否存在
    在尝试恢复 Flink 任务前,你可以先去 MySQL 服务器上确认 Flink 需要的 binlog 文件是否还存在。

    查看 Flink 需要的 binlog 位置: 这个信息通常在 Flink 的错误日志中可以找到,或者在 Flink UI 的 Checkpoint 详情里。

    查看 MySQL 服务器上可用的 binlog 文件: 在 MySQL 中执行:

    SQL

    SHOW BINARY LOGS;
    这个命令会列出所有当前可用的 binlog 文件。检查一下 Flink 需要的文件是否在这个列表里。如果不在,就说明它已经被清理了。

    解决方案 3:放弃旧状态,从最新位置启动(数据会丢失!)

    这是一种下策,只有在你可以接受数据丢失(从任务停止到重启这段时间的数据)的情况下才能使用。

    当你确定无法从旧的 savepoint 恢复时,你可以选择放弃这个 savepoint,然后重新启动一个新任务,并设置启动参数从最新的位置开始消费。

    Java

    // Flink SQL
    'scan.startup.mode' = 'latest-offset'

    // DataStream API
    MySqlSource.builder()
    .startupOptions(StartupOptions.latest())
    // ... 其他配置
    .build();
    ⚠️ 警告:这样做会导致停机期间的所有数据变更全部丢失,请务必谨慎操作!

    2025-10-30 16:34:11
    赞同 259 展开评论
  • 俺也一样

    2025-10-30 10:35:06
    赞同 265 展开评论
  • 那些看似波澜不惊的日复一日,总有一天会看到坚持的意义!

    哈哈,看不懂,完全是问了完成新手任务而来

    2025-10-28 16:46:03
    赞同 263 展开评论
  • 摸鱼来的

    真心看不懂,完全是为了完成新手任务而来

    2025-09-24 16:34:48
    赞同 254 展开评论
  • 大佬看不懂啊1

    2025-09-22 17:11:31
    赞同 201 展开评论
  • 真心看不懂,完全是为了完成新手任务而来

    2025-09-05 11:11:19
    赞同 198 展开评论
  • 一名技术小白正在学习中....

    真心看不懂,完全是为了完成新手任务而来

    2025-08-20 23:13:42
    赞同 196 展开评论
  • 哈哈,看不懂,完全是问了完成新手任务而来

    2025-07-25 07:58:21
    赞同 209 展开评论
滑动查看更多

实时计算Flink版是阿里云提供的全托管Serverless Flink云服务,基于 Apache Flink 构建的企业级、高性能实时大数据处理系统。提供全托管版 Flink 集群和引擎,提高作业开发运维效率。

还有其他疑问?
咨询AI助理