Apache Flink CancelJob

Apache Flink CancelJob

Certified

Cancel a running Flink job

Stops a Flink job via the REST API, optionally creating a savepoint first. Supports draining streaming jobs and waits until the cancellation completes or times out.

yaml
type: io.kestra.plugin.flink.CancelJob

Cancel a job with savepoint

yaml
id: cancel-flink-job
namespace: company.team

tasks:
  - id: cancel-job
    type: io.kestra.plugin.flink.CancelJob
    restUrl: "http://flink-jobmanager:8081"
    jobId: "{{ inputs.jobId }}"
    withSavepoint: true
    savepointDir: "s3://flink/savepoints/canceled/{{ execution.id }}"
    drainJob: true

Force cancel without savepoint

yaml
id: force-cancel-job
namespace: company.team

tasks:
  - id: force-cancel
    type: io.kestra.plugin.flink.CancelJob
    restUrl: "http://flink-jobmanager:8081"
    jobId: "{{ inputs.jobId }}"
    withSavepoint: false
Properties

Job ID

ID of the Flink job to cancel.

Flink REST API URL

Base URL of the Flink REST API (e.g., http://flink-jobmanager: 8081).

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.

Definitions
assetFailureBehaviorstring
Possible Values
IGNOREFAILWARN

Asset 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.

enableAutobooleanstring

Whether to auto-register assets referenced dynamically at runtime that are not statically declared in inputs or outputs.

inputsarray

The assets consumed as inputs.

id*string
Min length1
typestring
outputs

The assets produced as outputs.

id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*object
descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
id*string
Min length1
Max length150
type*string
Min length1

Custom asset type

descriptionstring
displayNamestring
metadataobject
Default{}
namespacestring
Min length1
Max length150
Default60

Cancellation timeout

Maximum time to wait for cancellation completion in seconds; defaults to 60.

Defaultfalse

Drain job

Drain streaming input before stopping. Applicable to streaming jobs only; defaults to false.

Savepoint directory

Target directory for the savepoint; required when withSavepoint is true.

Defaultfalse

Create savepoint before cancellation

Trigger a savepoint before cancelling the job; defaults to false.

Cancellation result

Result message returned by Flink for the cancellation request.

Cancelled job ID

ID of the Flink job that was cancelled.

Savepoint path

Path to the savepoint created before cancellation when requested.