公司动态
掌握Savepoint管理:flink-on-k8s-operator实现Flink作业状态持久化的完整指南
掌握Savepoint管理flink-on-k8s-operator实现Flink作业状态持久化的完整指南【免费下载链接】flink-on-k8s-operator[DEPRECATED] Kubernetes operator for managing the lifecycle of Apache Flink and Beam applications.项目地址: https://gitcode.com/gh_mirrors/fli/flink-on-k8s-operatorFlink-on-k8s-operator是一款强大的Kubernetes operator专为管理Apache Flink和Beam应用的生命周期而设计。本指南将详细介绍如何利用该operator实现Flink作业状态的持久化管理特别是通过savepoint功能确保数据流处理的连续性和可靠性。什么是Flink SavepointFlink Savepoint是流处理作业执行状态的一致性镜像允许用户捕获正在运行的作业状态并在后续从该状态恢复作业。这种机制对于版本升级、故障恢复和系统迁移至关重要。如何从Savepoint启动作业要从已有的savepoint启动Flink作业只需在作业规范中指定fromSavepoint属性。以下是一个示例配置apiVersion: flinkoperator.k8s.io/v1beta1 kind: FlinkCluster metadata: name: flinkjobcluster-sample spec: ... job: fromSavepoint: gs://my-bucket/savepoints/savepoint-123 allowNonRestoredState: false ...allowNonRestoredState属性控制是否允许非恢复状态更多信息请参考Flink CLI文档。四种创建Savepoint的方法1. 自动Savepoint通过设置autoSavepointSeconds和savepointsDir属性operator可以自动定期创建savepoint。以下配置将每300秒在GCS存储桶中创建一个savepointapiVersion: flinkoperator.k8s.io/v1beta1 kind: FlinkCluster metadata: name: flinkjobcluster-sample spec: ... job: autoSavepointSeconds: 300 savepointsDir: gs://my-bucket/savepoints/ ...您可以通过以下命令检查savepoint状态kubectl describe flinkclusters flinkjobcluster-sample成功的savepoint会增加作业状态中的savepoint generation并记录最新的savepoint位置。2. 通过更新FlinkCluster自定义资源手动触发savepoint的一种方法是编辑作业规范中的savepointGeneration将其设置为当前状态中的savepointGeneration 1然后应用更新后的YAML文件apiVersion: flinkoperator.k8s.io/v1beta1 kind: FlinkCluster metadata: name: flinkjobcluster-sample spec: ... job: savepointsDir: gs://my-bucket/savepoints/ savepointGeneration: 3 ...应用更新kubectl apply -f flinkjobcluster_sample.yaml3. 通过附加注解到FlinkCluster自定义资源您还可以通过为FlinkCluster的元数据附加控制注解来触发savepointmetadata: annotations: flinkclusters.flinkoperator.k8s.io/user-control: savepoint使用kubectl命令添加注解kubectl annotate flinkclusters flinkjobcluster-sample flinkclusters.flinkoperator.k8s.io/user-controlsavepoint完成后您可以在控制状态和作业状态中检查进度和结果。4. 使用Flink CLI或REST API在某些情况下例如未指定savepointsDir您可能需要绕过operator直接使用Flink CLI或REST API创建savepoint。首先创建到JobManager服务UI端口的端口转发kubectl port-forward svc/[FLINK_CLUSTER_NAME]-jobmanager 8081:8081然后使用Flink CLI创建savepointflink savepoint -m localhost:8081 [JOB_ID] [SAVEPOINT_DIR]或者调用Flink API触发异步savepoint操作curl -X POST -d {target-directory: [SAVEPOINT_DIR], cancel-job: false} http://localhost:8081/jobs/[JOB_ID]/savepoints从最新Savepoint自动重启作业长时间运行的作业可能因各种原因失败。通过将restartPolicy属性设置为FromSavepointOnFailureoperator可以在作业失败时自动从最新的savepoint重启作业apiVersion: flinkoperator.k8s.io/v1beta1 kind: FlinkCluster metadata: name: flinkjobcluster-sample spec: ... job: autoSavepointSeconds: 300 savepointsDir: gs://my-bucket/savepoints/ restartPolicy: FromSavepointOnFailure ...注意只有当作业状态中记录了savepoint时operator才能自动重启失败的作业作业状态中的fromSavepoint属性显示作业实际启动或重启的savepoint在远程存储中存储Savepoints通常您会希望将savepoints存储在远程存储中。有关如何在GCS中存储savepoints的详细信息请参阅存储文档。总结通过flink-on-k8s-operator您可以轻松实现Flink作业的savepoint管理包括自动创建、手动触发和从savepoint恢复等功能。这不仅提高了作业的可靠性还简化了版本升级和系统维护流程。无论您是Flink新手还是有经验的用户本指南都能帮助您充分利用savepoint功能确保流处理作业的稳定运行。要开始使用flink-on-k8s-operator请克隆仓库git clone https://gitcode.com/gh_mirrors/fli/flink-on-k8s-operator更多详细信息请参考完整的用户指南和Savepoint管理指南。【免费下载链接】flink-on-k8s-operator[DEPRECATED] Kubernetes operator for managing the lifecycle of Apache Flink and Beam applications.项目地址: https://gitcode.com/gh_mirrors/fli/flink-on-k8s-operator创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考