Repository navigation
Fix XCom loss in deferrable KubernetesPodOperator when pod is already terminal - #72821
Draft
shepherd44 wants to merge 2 commits into
Draft
shepherd44 wants to merge 2 commits into
shepherd44 wants to merge 2 commits into
Conversation
… terminal When the pod reaches a terminal state before the operator defers, `invoke_defer_method` calls `trigger_reentry` inline instead of deferring. `trigger_reentry` returns the XCom sidecar output, but that value was dropped: neither `invoke_defer_method`, nor `execute_async`, nor the deferrable branch of `execute` returned it. In the regular deferral path Airflow takes the return value of the resume method (`trigger_reentry`) and stores it as `return_value`, so the loss only happens on this shortcut. `execute_sync` already returns its result, which is why non-deferrable runs are unaffected. The task still succeeds and `pod_name` / `pod_namespace` are pushed by `execute_async`, so the failure is silent: `do_xcom_push=True` produces no `return_value` and downstream tasks pulling that XCom get `None`. Pods that finish within a second or two hit this often. Propagate the return value through the three call sites and widen the two return annotations accordingly.
|
Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contributors' Guide
|
`SparkKubernetesOperator.execute` dropped the deferrable result the same way, so the XCom fix would have stopped at `KubernetesPodOperator` and Spark tasks would still lose `return_value`. `test_execute_deferrable_does_not_call_super` asserted `result is None`, which pinned the dropped-return behaviour rather than the "does not call super" contract the test is named for. Updated to assert the result reaches the caller.
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Why
When the pod reaches a terminal state before the operator defers,
invoke_defer_methodcallstrigger_reentryinline instead of deferring:trigger_reentryends with:That value was dropped at every level above it —
invoke_defer_method,execute_async, and the deferrable branch ofexecuteall called down without returning.In the regular deferral path Airflow takes the return value of the resume method (
trigger_reentry) and stores it asreturn_value, so the loss only happens on this shortcut.execute_syncalready returns its result, which is why non-deferrable runs are unaffected.The failure is silent. The task still succeeds, and
pod_name/pod_namespaceare pushed separately byexecute_async, so the XCom entries look populated:pod_name,pod_namespacepod_name,pod_namespace,return_valueDownstream tasks pulling that XCom get
None. In our deployment a Jinja template doing{{ ti.xcom_pull(task_ids="fetch")["data_key"] }}failed with'None' has no attribute 'data_key', three retries deep, while the upstream task was reported assuccess. Retrying cannot help — the XCom is already gone.Task logs make the two paths easy to tell apart:
Pods that finish within a second or two hit this often, which is why it showed up intermittently across several DAGs rather than as a hard failure.
What
Propagate the return value through the call sites, and widen the two return annotations (
-> None→-> Any) accordingly.SparkKubernetesOperator.executedropped the deferrable result the same way, so it is fixed in the same commit — otherwise the fix would stop atKubernetesPodOperatorand Spark tasks would still losereturn_value.New tests:
test_invoke_defer_method_returns_trigger_reentry_result_when_pod_already_terminal— the inline call's result is returnedtest_execute_returns_deferrable_result—executehands the deferrable result back, as the synchronous branch already doesBoth fail on
mainwithassert None == {'key': 'value'}and pass with the change.One existing test was updated:
test_execute_deferrable_does_not_call_superassertedresult is None, which pinned the dropped-return behaviour rather than the "does not call super" contract the test is named for. It now asserts the result reaches the caller.Scope notes
invoke_defer_methodare unaffected.GKEStartPodOperator(google) only callsself.defer(...), so it has no shortcut path.EksPodOperator(amazon) has the same inline shortcut and still returnsNone; this change neither breaks it nor fixes it. I left it out to keep this PR to one provider — it looks worth a separate PR.container_logssubstring check in the synchronous branch), and Fix multiple_outputs no-op on deferrable KubernetesPodOperator #67226 fixed a related case where the deferrable path droppedmultiple_outputshandling.Testing
300 passed. Two failures and ten errors (
TestSuppress,write_logs) reproduce identically on unmodifiedmainon my machine, so they are local environment issues rather than regressions from this change.ruff checkandruff format --checkpass on all touched files.Gen-AI disclosure
This PR was prepared with the assistance of Gen-AI tools. The root cause was located by reading provider source and correlating it against Airflow API data and task logs from a real deployment where the bug occurred. I reviewed the diff and the tests, ran them locally, and I am able to explain and stand behind the change.