Flink Checkpoint停滞超时及相关延迟错误排查求助
Flink 1.18 K8s Operator部署下Checkpoint超时与GCS高延迟问题排查
问题现象
将原Flink 1.15作业升级为Flink Kubernetes Operator部署的1.18版本后,出现以下异常:
- 配置的10秒一次Checkpoint每次触发后都会停滞直至10分钟超时
- 必须等待上一次Checkpoint完全超时后,才会触发下一次Checkpoint
- Task Manager日志中频繁出现GCS操作高延迟告警(如
op_create、stream_write_close_operations),并伴随目录修复日志(repairImplicitDirectory)
关键配置与日志信息
核心配置片段
spec: flinkConfiguration: # Checkpoint配置 execution.checkpointing.interval: 10s execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.unaligned: "true" # 高可用配置 high-availability.storageDir: gs://.../my-job/ha # 自适应调度与自动扩缩容 job.scheduler: adaptive kubernetes.operator.job.autoscaler.enabled: "true" # 存储目录 state.checkpoints.dir: gs://.../my-job/checkpoints podTemplate: spec: containers: - env: - name: ENABLE_BUILT_IN_PLUGINS value: flink-gs-fs-hadoop-1.18.1.jar
典型高延迟日志
Apr 15, 2024 3:44:20 PM com.google.cloud.hadoop.fs.gcs.GhfsStorageStatistics updateMinMaxStats INFO: Detected potential high latency for operation op_create. latencyMs=695; previousMaxLatencyMs=0; operationCount=3; context=gs://.../my-job/checkpoints/93adc4ebde39a17486133f0 ... Apr 15, 2024 4:04:20 PM com.google.cloud.hadoop.gcsio.GoogleCloudStorageFileSystem repairImplicitDirectory INFO: Successfully repaired 'gs://.../my-job/checkpoints/93adc4ebde39a17486133f0c3ef3f508/chk-3855/' directory.
排查方向与验证步骤
1. GCS Hadoop插件行为变化
Flink 1.18的flink-gs-fs-hadoop插件依赖的GCS Hadoop客户端版本较1.15有更新,其中隐式目录修复功能可能默认开启且行为更激进:
- 日志中的
repairImplicitDirectory操作会在写入Checkpoint时自动修复目录结构,大量Checkpoint场景下会累积延迟 - 验证方式:在
flinkConfiguration中添加fs.gs.implicit.dir.repair.enabled: "false",重启作业后观察Checkpoint是否正常完成
2. 自适应调度器与自动扩缩容的冲突
启用adaptive调度器和K8s Operator自动扩缩容后,频繁的并行度调整会干扰Checkpoint流程:
- 未对齐Checkpoint(
unaligned: true)在并行度变化时会产生更多的GCS写入请求,加剧延迟 - 自动扩缩容的3分钟 metrics 窗口与10秒Checkpoint间隔不匹配,导致Checkpoint在集群波动期间无法正常结束
- 验证方式:临时禁用自动扩缩容(
kubernetes.operator.job.autoscaler.enabled: "false")并切换回默认调度器(job.scheduler: default),观察Checkpoint状态
3. Checkpoint间隔与存储能力不匹配
10秒的Checkpoint间隔过于频繁,结合GCS的天然延迟,会导致Checkpoint任务积压:
- 每个Checkpoint的写入时间超过间隔,后续Checkpoint被迫等待上一次超时,形成恶性循环
- 验证方式:临时将
execution.checkpointing.interval调整为30秒,同时缩短execution.checkpointing.timeout至3分钟,观察是否还会出现超时
4. GCS存储环境问题
- 确认K8s集群与GCS存储桶是否在同一区域,跨区域访问会大幅增加延迟
- 检查GCS监控指标(写入延迟、请求成功率),确认是否是存储端的性能瓶颈
- 验证凭证权限:确保Pod挂载的
GOOGLE_APPLICATION_CREDENTIALS拥有GCS的读写权限,且网络策略未限制GCS访问
优化建议
调整GCS客户端配置,减少IO操作延迟:
flinkConfiguration: fs.gs.implicit.dir.repair.enabled: "false" fs.gs.fast.upload: "true" # 启用并行上传 fs.gs.outputstream.sync.interval: "1048576" # 增大同步间隔,减少频繁IO优化Checkpoint与扩缩容的协同:
- 延长自动扩缩容的稳定间隔:
kubernetes.operator.job.autoscaler.stabilization.interval: "5m",避免频繁调整并行度 - 若必须使用自适应调度器,开启Checkpoint期间暂停扩缩容(Flink 1.18+支持通过
kubernetes.operator.job.autoscaler.checkpoint.pause: "true"配置)
- 延长自动扩缩容的稳定间隔:
调整Checkpoint参数适配存储能力:
- 将Checkpoint间隔调整为30秒~1分钟,根据GCS实际写入延迟确定合理值
- 缩短Checkpoint超时时间至3分钟,避免长时间阻塞作业
内容的提问来源于stack exchange,提问作者Rion Williams
相关产品推荐
相关产品推荐

