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

PyFlink 1.18在Kubernetes中重启后无法从Checkpoint恢复的解决方法

PyFlink作业重启后没法从Checkpoint恢复?看这几个修复点

问题根源

你的代码里有几个关键配置没到位,直接导致Pod重启后作业不从Checkpoint恢复,还丢Kafka数据:

  1. 硬设Kafka从最新偏移量开始消费:kafka_consumer.set_start_from_latest()这行代码会让作业每次启动都直接读最新消息,完全忽略Checkpoint里存的偏移量,等于白做Checkpoint。
  2. 没配置Checkpoint保留策略:Flink默认作业停了就删Checkpoint,重启的时候根本没可用的Checkpoint来恢复状态。
  3. 缺状态后端配置:虽然指定了S3路径,但没明确配置FileSystemStateBackend,状态可能没正确持久化到S3里。

一步步修复

1. 删掉Kafka消费者的set_start_from_latest()

开了Checkpoint又设了set_commit_offsets_on_checkpoints(True),Flink会自动从Checkpoint里恢复Kafka偏移量,根本不用手动指定起始位置。留着这行,重启就直接从最新位置开始,中间数据肯定丢。

改完的Kafka消费者代码:

# 创建Kafka消费者
kafka_consumer = FlinkKafkaConsumer(
    topics=[
        "topic-1", 
        "topic-2"
    ],
    deserialization_schema=SimpleStringSchema(),
    properties=kafka_consumer_properties
)
kafka_consumer.set_commit_offsets_on_checkpoints(True)
# 删掉 kafka_consumer.set_start_from_latest() 这行

2. 加上Checkpoint保留配置

得告诉Flink,作业停了别删Checkpoint,还要保留最近的几个,重启的时候才能用上:

# checkpoint 配置
env.enable_checkpointing(60*1000, CheckpointingMode.AT_LEAST_ONCE)
s3_checkpoint_path = "s3://xxx/flink-checkpoints"
# 明确配置FileSystemStateBackend
state_backend = FileSystemStateBackend(s3_checkpoint_path)
env.setStateBackend(state_backend)
# 配置Checkpoint保留策略
checkpoint_config = env.get_checkpoint_config()
checkpoint_config.set_checkpoint_storage(s3_checkpoint_path)
# 作业停止后保留Checkpoint,哪怕是被取消的情况
checkpoint_config.set_retention_checkpoint_after_job_stops(CheckpointRetentionPolicy.RETAIN_ON_CANCELLATION)
# 最多保留1个最新的Checkpoint,避免占太多S3空间
checkpoint_config.set_max_retained_checkpoints(1)

3. 配置作业重启策略

在K8s里,得让Flink作业在Pod挂了之后自动重启,不然得手动启动,这步不能少:

# 配置重启策略:最多重启3次,每次间隔10秒
env.set_restart_strategy(RestartStrategies.fixed_delay_restart(
    3,
    Time.seconds(10)
))

4. 验证S3的权限和路径

  • 确保Pod能访问S3桶:可以在Pod里跑s3cmd ls s3://xxx/flink-checkpoints看看能不能列出文件,要是权限不够,得给Pod绑IAM角色或者配置Access Key。
  • 检查S3路径对不对:看看路径下有没有chk-xxx开头的文件夹,那就是生成的Checkpoint文件,要是没有,说明Checkpoint根本没写进去。

5. 确认EventTime和Watermark配置没问题

你用的是for_monotonous_timestamps(),得保证自定义的CustomTimestampAssigner能正确提取事件时间,不然窗口触发不了,状态也存不到Checkpoint里,等于白搭。

修复后的完整代码示例

# 初始化流处理执行环境
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(1)
env.set_stream_time_characteristic(TimeCharacteristic.EventTime)

# 配置重启策略
env.set_restart_strategy(RestartStrategies.fixed_delay_restart(
    3,
    Time.seconds(10)
))

# checkpoint 相关配置
env.enable_checkpointing(60*1000, CheckpointingMode.AT_LEAST_ONCE)
s3_checkpoint_path = "s3://xxx/flink-checkpoints"
# 配置状态后端
state_backend = FileSystemStateBackend(s3_checkpoint_path)
env.setStateBackend(state_backend)
# 配置Checkpoint保留策略
checkpoint_config = env.get_checkpoint_config()
checkpoint_config.set_checkpoint_storage(s3_checkpoint_path)
checkpoint_config.set_retention_checkpoint_after_job_stops(CheckpointRetentionPolicy.RETAIN_ON_CANCELLATION)
checkpoint_config.set_max_retained_checkpoints(1)

# Kafka消费者配置
kafka_consumer_properties = {
    "bootstrap.servers": self.kafka_host,
    "group.id": "xxx_group",
    "security.protocol": "PLAINTEXT"
}

# 创建Kafka消费者
kafka_consumer = FlinkKafkaConsumer(
    topics=[
        "topic-1", 
        "topic-2"
    ],
    deserialization_schema=SimpleStringSchema(),
    properties=kafka_consumer_properties
)
kafka_consumer.set_commit_offsets_on_checkpoints(True)
# 已移除set_start_from_latest()

# Watermark策略配置
watermark_strategy = (
    WatermarkStrategy.for_monotonous_timestamps().with_timestamp_assigner(
        self.CustomTimestampAssigner()
    )
)

# 构建作业流
stream = env.add_source(kafka_consumer) \
    .map(self._decrypt_and_unpickle_message) \
    .filter(lambda x: x is not None) \
    .map(self.NormalizeMap(self)) \
    .assign_timestamps_and_watermarks(watermark_strategy) \
    .key_by(lambda record: record.get("_dwh_table_name"), key_type=Types.STRING()) \
    .window(TumblingEventTimeWindows.of(Time.minutes(1))) \
    .process(self.Sink(self))

stream.print()

# 执行作业
env.execute("Kafka Flink Streaming Job")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 21:07:01