Ian Stewart created FLINK-40918:
-----------------------------------
Summary: SIGTERM during application-mode HA recovery permanently
fails the suspended jobs of a PyFlink application
Key: FLINK-40918
URL: https://issues.apache.org/jira/browse/FLINK-40918
Project: Flink
Issue Type: Bug
Components: API / Python, Runtime / Coordination
Affects Versions: 2.3.0
Reporter: Ian Stewart
In Flink 2.3.0 application mode with HA, a graceful JobManager shutdown
(SIGTERM) that arrives after a new JobManager has recovered a job from HA, but
before the user \{{main()}} has resubmitted it, marks that job globally FAILED
with "Job recovery is not needed." Flink then clears the job's HA checkpoint
pointers, removes its execution plan and writes a JobResultStore entry. Every
later JobManager ignores the resubmission as globally terminal, so the job
never runs again. For a PyFlink application the window lasts from JobManager
start until the Python driver resubmits the job.
h3. Environment
* Flink 2.3.0, application mode on Kubernetes (Flink Kubernetes Operator
1.16.1), Kubernetes HA, \{{FileSystemJobResultStore}}.
* PyFlink streaming jobs: entry class
\{{org.apache.flink.client.python.PythonDriver}}, detached \{{execute()}},
checkpointing enabled.
h3. What happened
A node drain sent SIGTERM to a running JobManager. Its replacement recovered 1
job graph from HA, logged \{{Successfully recovered 0 persisted applications}},
started the Python driver (\{{Python Process Started}}), and about 10 s later
was itself sent SIGTERM by a second drain, before the driver had resubmitted
the job. Within about 100 ms the job was FAILED and its HA data removed. Two
unrelated jobs failed the same way in the same drain. The JobResultStore entry:
{code}
org.apache.flink.util.FlinkException: Job recovery is not needed.
at
org.apache.flink.runtime.dispatcher.Dispatcher.cleanupRemainingSuspendedJobs(Dispatcher.java:1077)
at
org.apache.flink.runtime.dispatcher.Dispatcher.notifyApplicationStatusChange(Dispatcher.java:1020)
at
org.apache.flink.runtime.application.AbstractApplication.lambda$transitionState$0(AbstractApplication.java:255)
at java.base/java.util.ArrayList.forEach(Unknown Source)
at
org.apache.flink.runtime.application.AbstractApplication.transitionState(AbstractApplication.java:253)
at
org.apache.flink.runtime.application.AbstractApplication.transitionToFailed(AbstractApplication.java:229)
at
org.apache.flink.client.deployment.application.PackagedProgramApplication.lambda$finishAsFailed$7(PackagedProgramApplication.java:510)
...
{code}
Every later JobManager logs \{{Successfully recovered 0 persisted job graphs}},
and the driver's resubmit is rejected with \{{Ignoring job submission ...
because the job already reached a globally-terminal state}}. The Kubernetes
operator reports the job as not found rather than FAILED, so
\{{kubernetes.operator.job.restart.failed}} does not apply. The only recovery
is to redeploy from a savepoint.
The INFO/WARN lines for the shutdown-time steps below were not captured in our
logs, although the stack trace shows that the code ran. The JobResultStore
entry is the reliable evidence.
h3. Root cause (release-2.3.0)
# Application mode parks jobs recovered from HA in
\{{Dispatcher#suspendedJobs}} (Dispatcher.java:398-410). A parked job resumes
only when the user \{{main()}} resubmits the same job id and the client calls
\{{Dispatcher#recoverJob}} (Dispatcher.java:885).
# \{{PythonDriver#main}} registers
\{{PythonEnvUtils.PythonProcessShutdownHook}} (PythonDriver.java:105-108). On
SIGTERM the JVM runs it at once and it destroys the Python process. The
driver's \{{main()}} then throws, and
\{{PackagedProgramApplication#onApplicationCanceledOrFailed}} takes the
"Application failed unexpectedly" branch
(PackagedProgramApplication.java:471-475) into \{{finishAsFailed}} and
\{{transitionToFailed}} (:510).
# The branch that would make this harmless is the \{{isDisposing}} cancellation
(PackagedProgramApplication.java:457-460). It needs
\{{PackagedProgramApplication#dispose}}, which is called only from
\{{Dispatcher#onStop}} through \{{terminateRunningApplications}}
(Dispatcher.java:805, 2243-2255). On SIGTERM,
\{{DispatcherResourceManagerComponent#internalShutdown}} first closes the
dispatcher operation caches and the REST endpoint, and stops the dispatcher
only after that (DispatcherResourceManagerComponent.java:145-159). The shutdown
hook wins this race.
# \{{transitionToFailed}} calls \{{Dispatcher#notifyApplicationStatusChange}}.
Its \{{cleanupRemainingSuspendedJobs}} (Dispatcher.java:1020, 1066-1087) builds
a FAILED \{{JobResult}} "Job recovery is not needed." for every job still
parked and runs the cleanup runner. The completed checkpoint store is shut down
as FAILED and cleared, the execution plan is removed, and the JobResultStore
entry is written.
h3. Regression from 2.2
In 2.2.1, recovered jobs started at once. An unexpected \{{main()}} failure was
handed to \{{FatalErrorHandler#onFatalError}}
(\{{ApplicationDispatcherBootstrap#finishBootstrapTasks}}), so the JobManager
failed over with HA data intact. In 2.3.0, any death of the user \{{main()}}
during the recovery window is permanent for the parked jobs. Failing parked
jobs when an application genuinely fails appears intended (FLINK-38974 /
FLIP-560). Treating a JVM shutdown as an application failure is the bug.
PyFlink is the exposed case because PythonDriver's shutdown hook turns SIGTERM
into a \{{main()}} failure. A Java \{{main()}} has no such hook, and \{{kill
-9}} runs no shutdown hooks at all. The 2.3 release testing of HA recovery
(FLINK-39525) used \{{kill -9}} with Java jobs, so neither path was exercised.
h3. Steps to reproduce
# Run a PyFlink streaming job in application mode with HA (ZooKeeper or
Kubernetes) and checkpointing enabled.
# Once the job is running, \{{kill -9}} the JobManager so that a new JobManager
recovers the job from HA.
# On the new JobManager, after \{{Python Process Started}} is logged and before
the job is resubmitted, send \{{kill -15}}.
# Start another JobManager.
Expected: the job is recovered from HA and resumes from its latest checkpoint.
Actual: \{{Successfully recovered 0 persisted job graphs}}. The JobResultStore
holds a FAILED entry with "Job recovery is not needed.", and the resubmission
is ignored as globally terminal.
h3. Possible fixes
* Do not fail suspended jobs when the application fails during the recovery
window. Either restore the 2.2 behaviour for that case (escalate to
\{{onFatalError}} so the JobManager fails over with HA intact), or clean up
remaining suspended jobs only when the application reaches FINISHED or CANCELED.
* Dispose running applications at the start of
\{{DispatcherResourceManagerComponent#internalShutdown}}, before the REST
endpoint closes, so that a \{{main()}} that dies during shutdown takes the
\{{isDisposing}} branch.
h3. Workaround
A Kubernetes \{{preStop}} hook on the JobManager that waits, for a bounded
time, while the PythonDriver's Python process is still running. This delays
SIGTERM until the job has been resubmitted. It does not cover a driver that
dies on its own inside the window.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)