Apache Flink从Checkpoint/Savepoint恢复状态的机制及懒加载疑问
How Apache Flink Restores State from Checkpoints/Savepoints
Let’s break down the core recovery process step by step, like you’re walking through it with a teammate:
- Kick off recovery: When you launch a job pointing to a Checkpoint or Savepoint (using
--fromSavepointor config settings), the JobManager first validates the path and pulls in the checkpoint metadata. This metadata has all the critical details: which operators hold state, where that state is stored, and the last processed offsets for every source. - Map state to tasks: The JobManager matches each state shard to the corresponding task/operator instance. Flink splits state by parallelism, so each TaskManager only grabs the state pieces assigned to its running tasks—no wasted bandwidth pulling irrelevant data.
- Spin up tasks and operators: TaskManagers start their assigned tasks and initialize operators. At this point, operators connect to their state backend (Heap, RocksDB, etc.) but don’t load any actual state data yet.
- Load state (when needed): Depending on your state backend setup, the actual state data is either loaded upfront or on demand. For most production setups (especially RocksDB), data is fetched from distributed storage (HDFS, S3, etc.) only when the operator needs to access it.
- Resume processing: Once all operators are ready and state is prepped, the job picks up exactly where the Checkpoint/Savepoint left off—using the recorded offsets to avoid data loss or duplication.
Is State Lazily Loaded from Checkpoints/Savepoints?
Your observation is totally on target, and yes—Flink uses lazy loading for state recovery, which is why start-to-first-event times are nearly identical for fresh jobs and restored ones. Here’s the breakdown:
- Metadata first, data later: During startup, only state metadata is loaded (things like state names, key group mappings, storage paths)—not the actual state values. This step is fast, no matter how big your state is.
- On-demand data fetch: The real state data is only pulled from storage when the operator first tries to use it (like when processing the first event that needs to read a keyed state entry). That means the time between job startup and receiving the first event isn’t affected by state size—since the heavy lifting of loading state happens after the first event arrives.
- Small caveat for huge state: While startup time stays the same, you might see a tiny delay when processing the first few events if your state is massive. That’s when Flink is pulling the necessary state shards into the backend—but this delay doesn’t show up in the "start to first event" metric, because it happens post-event-receipt.
So to confirm: what you’re seeing is exactly how Flink’s lazy state loading works.
内容的提问来源于stack exchange,提问作者Moein Hosseini
相关产品推荐
相关产品推荐

