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)

Reply via email to