Skip to content

[Fix-16789][Task-Plugin] Cancel Flink job through the flink CLI when stopping a task - #18667

Open
Admaing wants to merge 6 commits into
apache:devfrom
Admaing:fix/flink-16789-cancel-application
Open

Admaing wants to merge 6 commits into
apache:devfrom
Admaing:fix/flink-16789-cancel-application

Conversation

@Admaing

@Admaing Admaing commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Was this PR generated or assisted by AI?

YES. Prepared with AI assistance, using DSH (DeepSeek Harness) as the coding agent; the change was built, formatted and tested locally (mvn clean test + spotless:check on both Flink task modules, 24 tests) before submission.

Purpose of the pull request

Fix #16789. Stopping a Flink task killed only the client process tree, so the Flink job could keep running afterwards:

  • batch FlinkTask never issued flink cancel;
  • FlinkStreamTask did, but with an unresolved ${FLINK_HOME} — the cancel command is executed by ProcessBuilder, which does not expand shell variables — so it always failed with Cannot run program "${FLINK_HOME}/bin/flink".

Brief change log

  • resolve ${FLINK_HOME} for the cancel / savepoint commands (task params first, then the FLINK_HOME environment variable)
  • move the cancel logic from FlinkStreamTask into FlinkTask: cancel by the JobID printed by flink run, and fall back to the previous behaviour (kill the process tree, cancel the YARN/K8s application) when the task was submitted to a resource manager, when there is no JobID, or when the CLI fails
  • wait for the cancel command to finish instead of fire-and-forget, so a failure can be detected
  • Note: this changes the stop path of batch Flink tasks — they now attempt flink cancel before the process kill

Verify this pull request

This change added tests and can be verified as follows:

  • FlinkArgsUtilsTest: ${FLINK_HOME} is resolved in the cancel / savepoint command
  • FlinkTaskTest: cancel via the CLI using the JobID from the task log; fallback when the CLI fails; fallback when no JobID is found; no CLI call for YARN/K8s submissions
  • Manually verified with the real classes in a Linux container: a batch task (deployMode=cluster, detached submit) and a stream task both execute flink cancel <jobId> now, where before the batch task never issued it and the stream task always threw

Pull Request Notice

Pull Request Notice

…stopping a task

Stopping a Flink task only killed the client process tree, and the cancel path
of FlinkStreamTask executed `${FLINK_HOME}/bin/flink cancel <id>` with
${FLINK_HOME} left unresolved, because ProcessBuilder doesn't expand shell
variables. So the Flink job could keep running after the task was stopped: the
batch task never issued `flink cancel` at all, and the stream task always failed
with `Cannot run program "${FLINK_HOME}/bin/flink"`.

- resolve ${FLINK_HOME} for the cancel and savepoint commands (task params
  first, then the FLINK_HOME environment variable)
- move the cancel logic from FlinkStreamTask into FlinkTask, so batch and stream
  tasks share it: cancel by the JobID printed by `flink run` and fall back to the
  previous behaviour (kill the process tree and cancel the YARN/K8s application)
  when there is no JobID, when the task was submitted to a resource manager, or
  when the CLI fails
- wait for the cancel command to finish instead of fire-and-forget, so a failure
  can be detected and fall back
- add unit tests for the command resolution and the cancel/fallback logic

Verified by driving the real classes in a Linux container: a batch task
(deployment cluster, detached submit) and a stream task both execute
`flink cancel <jobId>` now, where previously the batch task never issued it and
the stream task always threw.

@SbloodyS SbloodyS left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Use the task environment for cancellation.

FlinkArgsUtils.buildFlinkCommand() only reads prepared parameters and System.getenv(), while submission also loads shell.env_source_list and TaskExecutionContext.environmentConfig. If FLINK_HOME is defined only in the selected task environment, submission succeeds but cancellation attempts the literal "${FLINK_HOME}/bin/flink" and fails. Please reuse the submission environment and tenant context.

Preserve cancellation failures and keep JobIDs separate from application IDs.

When the Flink CLI fails, FlinkTask.cancelApplication() falls back to killing the client and performing resource-manager cleanup. This cannot cancel a standalone/session Flink job. Moreover, cancelFlinkJob() has already overwritten context.appIds with the Flink JobID, so the fallback passes it to the YARN application manager. Please keep these identifiers separate and propagate the failure when the remote job cannot be cancelled.

Drain or redirect the CLI output.

FlinkArgsUtils.executeCommand() waits for completion without consuming stdout or stderr. If either pipe fills, the CLI blocks until the 30-second timeout and is forcibly terminated. Please drain or redirect both streams and retain the output for failure diagnostics.

…path

- run the cancel / savepoint command in the same environment as the task.
  The command is now executed by a shell interceptor which sources
  shell.env_source_list, applies the task's custom environment and runs as the
  task's tenant, so ${FLINK_HOME} is resolved even when it is only defined in
  the selected task environment. The Java-side placeholder resolution is dropped.
- keep the Flink JobID and the YARN/K8s application id separate. The JobID is no
  longer written into taskExecutionContext.appIds, so the fallback cannot pass it
  to the YARN application manager, and a failed cancel is now propagated instead
  of being hidden behind the fallback that cannot cancel a standalone/session job.
- consume the output of the cancel command so it cannot block on a full pipe, and
  keep the output for failure diagnostics.
- add a test that executes the cancel command with FLINK_HOME defined only in the
  task environment and writes more output than a pipe buffer holds.
@Admaing

Admaing commented Sep 30, 2026

Copy link
Copy Markdown
Contributor Author

Use the task environment for cancellation.

FlinkArgsUtils.buildFlinkCommand() only reads prepared parameters and System.getenv(), while submission also loads shell.env_source_list and TaskExecutionContext.environmentConfig. If FLINK_HOME is defined only in the selected task environment, submission succeeds but cancellation attempts the literal "${FLINK_HOME}/bin/flink" and fails. Please reuse the submission environment and tenant context.

Preserve cancellation failures and keep JobIDs separate from application IDs.

When the Flink CLI fails, FlinkTask.cancelApplication() falls back to killing the client and performing resource-manager cleanup. This cannot cancel a standalone/session Flink job. Moreover, cancelFlinkJob() has already overwritten context.appIds with the Flink JobID, so the fallback passes it to the YARN application manager. Please keep these identifiers separate and propagate the failure when the remote job cannot be cancelled.

Drain or redirect the CLI output.

FlinkArgsUtils.executeCommand() waits for completion without consuming stdout or stderr. If either pipe fills, the CLI blocks until the 30-second timeout and is forcibly terminated. Please drain or redirect both streams and retain the output for failure diagnostics.

@SbloodyS All three points are fixed in 691373a:

  1. Cancel now runs in the same environment as the task (shell.env_source_list + environmentConfig + tenant), so ${FLINK_HOME} is resolved at runtime.
  2. JobID and appId are kept separate, and a failed flink cancel now throws instead of falling back silently.
  3. CLI output is drained on a separate thread, so it cannot block, and the output is kept for diagnostics.
    Tests: flink 15 passed , flink-stream 9 passed; the Linux-only env test also passes in a Linux container. Please take another look.

@SbloodyS SbloodyS left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Propagate savepoint command failures.

In FlinkStreamTask.java:77, savePoint() ignores the boolean returned by FlinkArgsUtils.executeCommand(). That method returns false on startup errors, non-zero exit codes, timeouts, and interruption. Consequently, savePoint() returns normally and StreamingTaskInstanceOperatorImpl.triggerSavepoint() reports success even when the savepoint command failed.

Please throw TaskException when executeCommand() returns false, and cover the failure response in the savepoint tests.

…ream task

savePoint() ignored the result of the savepoint command, so
StreamingTaskInstanceOperatorImpl reported success even when the command failed
to start, exited with a non-zero code, timed out or was interrupted. It also
returned normally when no application id could be collected, without triggering
any savepoint at all.

- throw TaskException when the savepoint command fails
- throw TaskException when no application id can be found, instead of reporting a
  successful savepoint that never ran
- cover the savepoint command, the failed command and the missing application id
  in FlinkStreamTaskTest
@Admaing

Admaing commented Sep 30, 2026

Copy link
Copy Markdown
Contributor Author

@SbloodyS Fixed in e5abd15: savePoint() now throws TaskException when the savepoint command returns false, and also when the application id is missing (it used to return silently and report success). Failure cases are covered in FlinkStreamTaskTest. Full run: flink 16 tests and flink-stream 12 tests pass on Linux, spotless passes.

@SbloodyS SbloodyS left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Correct the JobID used by the new savepoint success test.

FlinkStreamTaskTest.testSavePointCommand() expects:
flink savepoint application_1700000000000_0001

However, flink savepoint requires a Flink JobID; the YARN application ID is a separate parameter. Mocking executeFlinkCommand() to return true makes this invalid command appear successful.

This identifier problem predates the PR, but the new test now locks in the incorrect behavior. Please use a real Flink JobID in the success test and obtain it through getFlinkJobIds(), keeping any YARN application ID separate.

`flink savepoint` takes the Flink JobID, not the YARN/K8s application id.
FlinkStreamTask.savePoint() used the application id returned by
getApplicationIds(), which cannot work for standalone/session deployments, and
the new savepoint test locked in a command that could never succeed.

- buildSavePointCommandLine() now takes the JobID explicitly
- savePoint() resolves the JobID through getFlinkJobIds() and no longer reads or
  writes taskExecutionContext.appIds, so the two identifiers stay separate
- tests use a real JobID, cover the case where only a YARN application id is
  present, and assert the application id is never passed to `flink savepoint`
…nd test

testExecuteCommandReuseTaskEnvironment does not depend on Linux specific tooling,
it only runs a bash script through the shell interceptor. The whole
FlinkArgsUtilsTest passes on macOS as well (9/9), so the condition can be dropped
and the test also runs for developers on macOS.
@Admaing

Admaing commented Sep 30, 2026

Copy link
Copy Markdown
Contributor Author

@SbloodyS Fixed in 1f35037: savePoint() now uses the Flink JobID from getFlinkJobIds() and passes it to buildSavePointCommandLine(jobId); the YARN application id is no longer used or written back. The success test uses a real JobID (with the application id present in the log to prove it is ignored), and the missing-id test asserts no CLI call and a TaskException. flink 16 + flink-stream 12 tests pass on Linux, spotless passes.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[Bug] [Flink] stop flink job error

2 participants