You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用StateBackend的Flink应用是否需要配置Checkpoint?

核心问题

  1. StateBackend是否需要Checkpoint?
  2. 配置了StateBackend的Flink应用是否需要启用Checkpoint?

解答

  • StateBackend本身不需要Checkpoint:StateBackend是Flink负责存储作业运行时状态的组件,它的作用是管理状态的读写与存储介质(内存、RocksDB、文件系统等)。Checkpoint是Flink实现故障容错的机制,用于生成状态快照,StateBackend只是Checkpoint快照的存储载体之一,两者属于不同层面的概念。
  • 配置StateBackend后是否启用Checkpoint看业务需求:
    • 若作业需要故障恢复能力(重启后从断点继续处理,保证Exactly-Once或At-Least-Once语义),必须启用Checkpoint,此时StateBackend会将Checkpoint快照持久化到指定存储位置。
    • 若作业是无状态的,或不需要故障恢复(如测试作业、一次性离线计算),可以不启用Checkpoint,但StateBackend仍会管理作业运行时的内存状态(比如MemoryStateBackend的内存存储)。

示例代码(Kotlin)

class EnrichmentStream {
    private val checkpointsDir  = "file://${System.getProperty("user.dir")}/checkpoints/"
    private val rocksDBStateDir = "file://${System.getProperty("user.dir")}/state/rocksdb/"

    companion object {
        @JvmStatic
        fun main(args: Array<String>) {
            EnrichmentStream().runStream()
        }
    }

    fun runStream() {
        val environment = StreamExecutionEnvironment
            .createLocalEnvironmentWithWebUI(Configuration())

        environment.parallelism = 3

        // Checkpoint配置
        environment.enableCheckpointing(5000)
        environment.checkpointConfig.minPauseBetweenCheckpoints = 100
        environment.checkpointConfig.setCheckpointStorage(checkpointsDir)

        val stateBackend = EmbeddedRocksDBStateBackend()
        stateBackend.setDbStoragePath(rocksDBStateDir)
        environment.stateBackend = stateBackend

        environment.checkpointConfig.externalizedCheckpointCleanup =
            CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION

        // 配置重启策略
        environment.restartStrategy = RestartStrategies.fixedDelayRestart(5, Time.seconds(5))

        val tableEnvironment = StreamTableEnvironment.create(environment)

        // 执行SQL查询查看目录、数据库和表
        tableEnvironment
            .executeSql("SHOW CATALOGS")
            .print()

        tableEnvironment
            .executeSql("SHOW DATABASES")
            .print()

        tableEnvironment
            .executeSql("SHOW TABLES")
            .print()

        tableEnvironment
            .executeSql(Queries.CREATE_SENSORS_TABLE)
            .print()

        tableEnvironment
            .executeSql(Queries.CREATE_READINGS_TABLE)
            .print()

        tableEnvironment
            .executeSql("SHOW TABLES")
            .print()

        tableEnvironment
            .executeSql(Queries.JOIN_SENSOR_READINGS_WITH_INFO_QUERY)
            .print()
    }
}

内容的提问来源于stack exchange,提问作者overexchange

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 18:34:54