
Apache Flink TriggerSavepoint
CertifiedCreate a savepoint for a running Flink job
Apache Flink TriggerSavepoint
Create a savepoint for a running Flink job
Requests a savepoint via the Flink REST API and waits for completion. Supports optional target directory, format type, and post-savepoint cancellation.
type: io.kestra.plugin.flink.TriggerSavepointExamples
Trigger savepoint with specific directory
id: create-savepoint
namespace: company.team
tasks:
- id: trigger-savepoint
type: io.kestra.plugin.flink.TriggerSavepoint
restUrl: "http://flink-jobmanager:8081"
jobId: "{{ inputs.jobId }}"
targetDirectory: "s3://flink/savepoints/backup/{{ execution.id }}"
savepointTimeout: 300
Trigger savepoint with default directory
id: default-savepoint
namespace: company.team
tasks:
- id: trigger-savepoint
type: io.kestra.plugin.flink.TriggerSavepoint
restUrl: "http://flink-jobmanager:8081"
jobId: "{{ inputs.jobId }}"
Properties
jobId *string
Job ID
ID of the Flink job to snapshot.
restUrl *string
Flink REST API URL
Base URL of the Flink REST API (e.g., http://flink-jobmanager: 8081).
assets
Assets this task consumes as inputs or produces as outputs, for lineage tracking and the asset graph (Enterprise Edition). A flow declaring this property on a task is rejected in the open-source edition.
io.kestra.core.models.assets.AssetsDeclaration
IGNOREFAILWARNAsset failure behavior
Behavior applied to the task state when a declared asset fails to render, emit, or be persisted (e.g. a lock conflict): FAIL escalates it to FAILED, WARN (default) warns it if it would otherwise succeed, IGNORE leaves the state untouched.
Whether to auto-register assets referenced dynamically at runtime that are not statically declared in inputs or outputs.
The assets consumed as inputs.
io.kestra.core.models.assets.AssetIdentifier
1The assets produced as outputs.
io.kestra.plugin.ee.assets.Dataset
1150{}1150io.kestra.plugin.ee.assets.File
1150{}1150io.kestra.plugin.ee.assets.Table
1150{}1150io.kestra.plugin.ee.assets.VM
1150{}1150io.kestra.core.models.assets.External
1150{}1150io.kestra.core.models.assets.Custom
11501Custom asset type
{}1150cancelJob booleanstring
falseCancel job after savepoint
Cancel the job after the savepoint completes; defaults to false.
formatType string
CANONICALFormat type
Savepoint format type: CANONICAL or NATIVE. Defaults to CANONICAL for broader compatibility.
savepointTimeout integerstring
300Savepoint timeout
Maximum wait time for savepoint completion in seconds; defaults to 300.
targetDirectory string
Target directory
Target directory for the savepoint; falls back to the cluster default when omitted.
Outputs
jobId string
Job ID
ID of the Flink job for which the savepoint was created.
requestId string
Request ID
Savepoint request ID used for status polling.
savepointPath string
Savepoint path
Path returned by Flink for the created savepoint.