优化Flink S3 Checkpoint性能:解决限流与触发延迟问题
Flink S3 Checkpoint限流与触发延迟问题排查与解决
问题背景
我们的Flink作业采用S3作为Checkpoint后端,开启增量Checkpoint,状态量约100-200GB;作业总并行度5000,包含三个算子:Kinesis Consumer(750并行度)、状态机算子(3500并行度)、S3 Sink(450并行度)。其中Checkpoint使用s3-presto库,S3 Sink使用s3-hadoop库。
最初遇到大量S3套接字超时和限流异常,添加以下配置后超时问题解决:
state.backend: filesystem state.checkpoints.dir: s3://prod-flink-test-checkpoint/_entropy_/flink/checkpoints state.backend.incremental: true s3.connection.maximum: 1000 s3.entropy.key: _entropy_ s3.entropy.length: 8 state.storage.fs.memory-threshold: 1000000 state.backend.fs.write-buffer-size: 4194304 presto.s3.socket-timeout: 1m presto.s3.connect-timeout: 1m presto.s3.max-connections: 6000
但目前仍存在S3限流(SlowDown异常),且Checkpoint触发严重延迟:配置的触发间隔为70秒,实际10-40分钟才触发一次。当前Checkpoint核心配置如下:
checkpoint-interval : 70000 min-pause-between-checkpoint : 15000 max-concurrent-checkpoint : 1 checkpoint-type: aligned Aligned-checkpoint-timeout : 30sec TolerableCheckpointFailureNumber :2 Checkpoint-mode : atleast-once
作业无明显背压,JobManager日志也无其他异常,限流异常栈信息如下:
Caused by: org.apache.flink.util.SerializedThrowable: java.io.IOException: Could not flush to file and close the file system output stream to s3://prod-flink-test-checkpoint/DHpRvkE8/flink/checkpoints/6ddf92e0c7902eba815543fca2fa2730/chk-6/d78e8ca4-2ec8-4732-86f2-7a3b85df2637 in order to obtain the stream state handle at org.apache.flink.runtime.state.filesystem.FsCheckpointStreamFactory$FsCheckpointStateOutputStream.closeAndGetHandle(FsCheckpointStreamFactory.java:425) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.runtime.state.DuplicatingCheckpointOutputStream.closeAndGetPrimaryHandle(DuplicatingCheckpointOutputStream.java:260) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.runtime.state.CheckpointStreamWithResultProvider$PrimaryAndSecondaryStream.closeAndFinalizeCheckpointStreamResult(CheckpointStreamWithResultProvider.java:113) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.runtime.state.heap.HeapSnapshotStrategy.lambda$asyncSnapshot$3(HeapSnapshotStrategy.java:181) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.runtime.state.SnapshotStrategyRunner$1.callInternal(SnapshotStrategyRunner.java:91) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.runtime.state.SnapshotStrategyRunner$1.callInternal(SnapshotStrategyRunner.java:88) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.runtime.state.AsyncSnapshotCallable.call(AsyncSnapshotCallable.java:78) ~[flink-dist-1.16.0.jar:1.16.0] at java.util.concurrent.FutureTask.run(FutureTask.java:266) ~[?:1.8.0_372] at org.apache.flink.util.concurrent.FutureUtils.runIfNotDoneAndGet(FutureUtils.java:540) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.streaming.api.operators.OperatorSnapshotFinalizer.<init>(OperatorSnapshotFinalizer.java:54) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.streaming.runtime.tasks.AsyncCheckpointRunnable.finalizeNonFinishedSnapshots(AsyncCheckpointRunnable.java:191) ~[flink-dist-1.16.0.jar:1.16.0] at org.apache.flink.streaming.runtime.tasks.AsyncCheckpointRunnable.run(AsyncCheckpointRunnable.java:124) ~[flink-dist-1.16.0.jar:1.16.0] ... 3 more Caused by: org.apache.flink.util.SerializedThrowable: java.io.IOException: com.amazonaws.services.s3.model.AmazonS3Exception: Please reduce your request rate. (Service: null; Status Code: 0; Error Code: SlowDown; Request ID: P947AG455MJ7TQB0; S3 Extended Request ID: VquhT/jyfuuLLjuJiFfrCCrXSMPw8fmkHGavMJfvMV/n2ENSENtQrhTJZ1B15n8s3PNDu5fsbAI=; Proxy: null), S3 Extended Request ID: VquhT/jyfuuLLjuJiFfrCCrXSMPw8fmkHGavMJfvMV/n2ENSENtQrhTJZ1B15n8s3PNDu5fsbAI=
问题根源拆解
- S3请求量远超配额:S3单桶默认每秒允许3500次PUT/COPY/POST/DELETE请求、5500次GET/HEAD请求。我们的作业有3500个状态机算子,加上增量Checkpoint的元数据读写,瞬间请求量直接打满S3配额,触发
SlowDown限流。 - 对齐Checkpoint的阻塞效应:配置了
checkpoint-type: aligned,对齐Checkpoint要求所有算子完成快照后才能结束整个Checkpoint流程。一旦某个算子因S3限流导致快照超时,整个Checkpoint会被卡住;再加上max-concurrent-checkpoint:1,前一个Checkpoint未完成,下一个无法触发,最终导致触发间隔被拉得极长。 - Checkpoint超时配置过短:当前
Aligned-checkpoint-timeout:30sec远不足以应对S3限流时的快照延迟,快照频繁超时失败,触发重试后进一步加剧S3请求压力,形成恶性循环。 - 增量Checkpoint元数据开销:大状态下的增量Checkpoint需要频繁读写manifest等元数据小文件,这类请求会额外占用S3配额,更容易触发限流。
排查思路
- 核对S3桶请求指标:查看CloudWatch中该S3桶的
NumberOfRequests和429Errors指标,确认是否达到S3请求配额,同时定位是PUT还是GET请求触发的限流。 - 分析Checkpoint阶段耗时:在Flink UI的Checkpoints页面,查看每个Checkpoint的触发、算子快照、完成等阶段的耗时,定位哪个算子的快照耗时最长,是否和S3限流时间点对应。
- 检查TaskManager资源:查看TaskManager的CPU、内存、网络IO使用率,确认是否存在资源瓶颈导致请求发送缓慢;同时查看TM日志中连接池相关日志,确认是否存在连接耗尽的情况。
- 验证配置生效情况:检查Flink日志,确认
presto.s3.max-connections等配置是否被正确加载,是否存在s3和presto.s3配置冲突的情况。
针对性优化方案
1. 缓解S3限流
- 减少小文件请求:调大
state.backend.fs.write-buffer-size至16MB或32MB,减少单次上传的请求次数;同时调大state.storage.fs.memory-threshold,让更多状态先在内存中缓冲,降低写S3的频率。 - 配置请求重试与限流:添加presto-s3的重试和限流参数,避免频繁重试加剧拥堵:
presto.s3.max-error-retries: 10 presto.s3.retry-base-delay: 100ms presto.s3.retry-max-delay: 5s presto.s3.throttle.enabled: true presto.s3.throttle.max-requests-per-second: 3000 # 略低于S3单桶PUT配额 - 分散请求压力:可以考虑将Checkpoint分散到多个S3桶,或者按算子/TaskManager划分目录前缀,进一步分散请求量。
2. 调整Checkpoint策略
- 切换为非对齐Checkpoint:对于大并行度、大状态的作业,启用
checkpoint-type: unaligned,允许算子在快照过程中继续处理数据,避免因单个算子快照阻塞整个Checkpoint流程,大幅缩短整体耗时。同时调整超时时间:checkpoint-type: unaligned unaligned-checkpoint.timeout: 10min - 调整Checkpoint并发与间隔:适当调大
max-concurrent-checkpoint至2(需确保集群资源足够),同时调整min-pause-between-checkpoint,避免过度限制触发频率。 - 延长对齐Checkpoint超时(若不切换非对齐):如果必须使用对齐Checkpoint,将
Aligned-checkpoint-timeout调整为5-10分钟,避免快照因限流频繁超时失败:aligned-checkpoint-timeout: 5min
3. 资源与连接池优化
- 分离Checkpoint与Sink的S3客户端:Checkpoint用s3-presto,Sink用s3-hadoop,两者连接池独立配置,避免资源抢占。比如给hadoop客户端单独设置连接池:
fs.s3a.connection.maximum: 1000 - 增加TaskManager网络资源:确保TaskManager有足够的网络带宽处理S3请求,避免网络瓶颈加剧延迟和限流。
4. 增量Checkpoint优化
- 启用旧Checkpoint清理:开启
state.backend.cleanup-incremental-checkpoint,自动清理旧的增量Checkpoint文件,减少S3文件数量和元数据请求量:state.backend.cleanup-incremental-checkpoint: true - 开启状态快照压缩:对于HeapStateBackend,启用压缩减少上传数据量:
state.backend.heap.snapshot-compression.enabled: true
验证方式
- 修改配置后,观察Flink UI中Checkpoint的触发间隔是否恢复到70秒左右。
- 查看CloudWatch中S3的429错误是否减少或消失。
- 查看Checkpoint各阶段耗时是否明显降低。
内容的提问来源于stack exchange,提问作者priyadhingra19
相关产品推荐
相关产品推荐

