PyFlink 1.18在Kubernetes中重启后无法从Checkpoint恢复的解决方法
PyFlink作业重启后没法从Checkpoint恢复?看这几个修复点
问题根源
你的代码里有几个关键配置没到位,直接导致Pod重启后作业不从Checkpoint恢复,还丢Kafka数据:
- 硬设Kafka从最新偏移量开始消费:
kafka_consumer.set_start_from_latest()这行代码会让作业每次启动都直接读最新消息,完全忽略Checkpoint里存的偏移量,等于白做Checkpoint。 - 没配置Checkpoint保留策略:Flink默认作业停了就删Checkpoint,重启的时候根本没可用的Checkpoint来恢复状态。
- 缺状态后端配置:虽然指定了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
相关产品推荐
相关产品推荐

