[Translation] Flink - Savepoint vs Checkpoint

Liao Jiayi Liao Jiayi #Flink#Checkpoint#Apache Flink

Translated from the Data Artisans blog, this article compares Savepoints and Checkpoints in Apache Flink.

Translated from Chinese with AI · Read the original

Translated from the dataArtisans blog: 3 differences between Savepoints and Checkpoints in Apache Flink. Many Flink developers confuse these two concepts. How do these seemingly similar features differ?

An Apache Flink Savepoint captures a snapshot of a running streaming program. It records the entire program state, including processing positions such as Kafka offsets. Flink uses the Chandy-Lamport snapshot algorithm to create consistent Savepoints. A Savepoint has two main elements:

  1. A binary file, usually large, recording all state in the streaming program.
  2. A relatively small metadata file, stored in the distributed filesystem or data store you specify, containing pointers (paths) to every Savepoint file.

For details, see our earlier step-by-step guide on how to enable Savepoints.

This sounds much like the Checkpoints discussed in earlier posts. Checkpointing is Flink’s internal mechanism for fault recovery, copying and persisting state, and tracking consumption positions. If a program fails, Flink restores checkpointed state and resumes processing from the pre-failure position, as though nothing happened.

See How Apache Flink manages Kafka Consumer offsets.

flink-savepoint-3

Checkpoints and Savepoints are distinctive features of Flink’s streaming framework. Their implementations are similar, but they differ in three ways:

  1. Purpose: Conceptually, Savepoints and Checkpoints differ like backups and recovery logs in traditional databases. Checkpoints enable recovery from potential failures, such as transient network problems. Savepoints let users explicitly trigger a backup and restore a program by restarting it.
  2. Implementation: Checkpoints are lightweight and fast, exploiting features of the underlying state store for quick backup and recovery. With RocksDB, for example, state is persisted in RocksDB’s format rather than Flink’s native format, enabling incremental Checkpoints. This accelerates checkpointing and was its first lighter-weight implementation. Savepoints emphasize portability and support arbitrary job changes, at the cost of more expensive backup and recovery.
  3. Lifecycle: Checkpoints are triggered automatically on a schedule. Flink creates, maintains, and deletes them without user intervention. Savepoints must be triggered, deleted, and managed by the user.
Dimension Checkpoints Savepoints
Purpose Recovery/failover after job failure Manual backup/restart/job recovery
Implementation Lightweight and fast Portable but more expensive
Lifecycle Controlled by Flink Controlled manually by the user

When Should Streaming Programs Use Savepoints?

Although streaming programs process unbounded data, they may need to reprocess data already consumed. Savepoints help with these scenarios:

  • Update production programs with new features, bug fixes, or better machine-learning models
  • Introduce A/B tests, comparing versions from the same point in the same source
  • Scale out to use more cluster resources
  • Use a new Flink version or move to a cluster running a newer version.

Conclusion

Checkpoints and Savepoints are different features, both supporting Flink’s consistency and fault tolerance. Checkpoints address possible program failures; Savepoints support upgrades, bug fixes, migration, and A/B testing. Together, they persist and restore program state across different scenarios.

Notes:

  1. For convenience, I replaced an image in the original article with a table.
  2. Some passages felt abstract when translated literally, so I added a few thoughts of my own.