Skip to content

Fetch several XCom values of a Taskinstance in one Execution API request - #70223

Draft
dabla wants to merge 43 commits into
apache:mainfrom
dabla:feature/add-xcoms-keys-route
Draft

dabla wants to merge 43 commits into
apache:mainfrom
dabla:feature/add-xcoms-keys-route

Conversation

@dabla

@dabla dabla commented Jul 22, 2026 •

Copy link
Copy Markdown
Contributor

Review only this change: the PR's own commit · diff against #62922's branch. Everything else in this PR's commit list and file list is #62922, which this PR is stacked on.

Stacked on #62922 (Task Iteration), which owns XComIterable: only the last commit, 76625b4, belongs to this PR. Squashed from this branch's earlier history and re-based on that PR, with the batched reads moved to XComIterable's current home in bases/xcom.py.

Add a POST /xcoms/{dag_id}/{run_id}/{task_id}/keys execution API endpoint that
accepts a list of XCom keys and returns all matching values in a single database
query. Use it in XComIterable to reduce iteration from N round-trips to one.

Motivation

XComIterable (introduced in AIP-104 Task Iteration) stores per-index results
under distinct keys (return_value_0, return_value_1, …) with the same
map_index. The existing GET …/slice endpoint cannot be reused because it
ranges over map_index for a single key — the inverse structure. Iterating or
slicing a 1000-item result previously issued 1000 separate XCom.get_one calls.

Changes

  • New endpoint POST /xcoms/{dag_id}/{run_id}/{task_id}/keys — accepts
    {"keys": [...]}, returns values in request-key order (null for missing
    keys), filtered by map_index query param (default -1)
  • has_xcom_access simplified: reads the optional key from
    request.path_params instead of requiring a {key} path segment, so a
    separate router/dependency is no longer needed
  • XComKeysRequest Pydantic model added to execution API data models
  • GetXComByKeys comms message, XComOperations.get_by_keys client method,
    handle_get_xcom_by_keys supervisor handler, and processor dispatch wiring
  • BaseXCom.get_by_keys sends the message and deserializes the values;
    XComIterable.__iter__ and __getitem__ slice read through it, one
    round-trip for a full iteration or slice. A single index stays one XCom read.
  • AddXComKeysEndpoint execution API version change and AddGetXComByKeys
    supervisor schema version entry, with the regenerated schema snapshot.
  • Tests for the new endpoint (TestGetXComByKeys) and for XComIterable
    (in test_xcom.py), plus supervisor message coverage.

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude Code following the guidelines

🤖 Generated with Claude Code

@amoghrajesh amoghrajesh left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

I haven't followed up closely with AIP 104, but I have some qns from an initial look.

Comment thread task-sdk/src/airflow/sdk/execution_time/schema/schema.json
Comment thread task-sdk/src/airflow/sdk/execution_time/lazy_sequence.py
Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py
Comment thread airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py
@dabla dabla added this to the Airflow 3.4.0 milestone Jul 22, 2026
@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch from 5258e09 to dbf25d5 Compare July 22, 2026 14:10
@dabla
dabla marked this pull request as draft August 17, 2026 08:52
@dabla

dabla commented Aug 17, 2026

Copy link
Copy Markdown
Contributor Author

This PR should only be merged once AIP-104 is merged.

@dabla

dabla commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

Interaction with #73807 to handle when both land: FlattenedXComIterable there inherits XComIterable.__iter__, which this PR turns into a single batched fetch of the raw pages. Whichever of the two merges second must give FlattenedXComIterable its own __iter__ over flattened positions; see the matching note on #73807.

@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch 12 times, most recently from 04cc40e to b74f42c Compare October 1, 2026 10:19
@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch from f2a2cbd to 5594e71 Compare October 9, 2026 14:18
The success and skip callbacks fired from IndexedTaskRunner.__exit__, as
soon as the operator returned and before the indexed task's checkpoint and
result were written. A checkpoint write that failed had already announced
a success, and the retry ran the indexed task again and announced it a
second time; the callback also found no end_date on the task instance. A
plain task fires its callback after the push and the state report, with
end_date set.

The exit now only notes a failure. IterableOperator._run_task calls the
runner's report_success() or report_skip() once the matching checkpoint is
written, where the indexed task's code ran (its worker thread for a sync
operator, the loop for an async one); both set end_date and the state
before the callback. A result push that fails after the checkpoint is
replayed from it, so the callback fires once. The docs page and the class
docstring say so.

test_the_success_callback_waits_for_the_checkpoint and the two runner
report tests fail before.
@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch from 5594e71 to f054261 Compare October 9, 2026 14:48
dabla added 2 commits October 9, 2026 17:50
…valuating the policy again

With several failed indexed tasks whose retry policy decisions were all the
default, or whose evaluation raised, no exception outweighed the others, the
group was handed to the runner and no decision was kept for it. The callbacks
then evaluated the policy once more, on the BaseExceptionGroup rather than on
any indexed task's exception: K+1 evaluations instead of K, and for a policy
that answers differently per call a retry announced here that the runner's own
evaluation of the group could refuse. The once-per-item test did not reach this
path because its policy answers retry(), so an exception was always chosen.

_failure_for_the_runner now keeps the default decision for the group when no
decision outweighs it, so _task_will_retry follows it and falls back to retry
eligibility without evaluating the policy again.

test_an_undecided_group_is_not_evaluated_again_for_the_callbacks fails before
with four evaluations.
…the iteration

on_kill() reaches the sub-operators that have started, and the stop flag keeps
the executor from pulling more. An indexed task instance is neither until its
runner starts: on a retry attempt it first reads its checkpoint, and a kill that
lands during that read found it unregistered, so it started afterwards, ran to
completion unkilled, its remote work included, and the drain waited for it
before the task could conclude as terminated.

IterableOperator._run_task now asks the iteration state whether a stop was
requested once the checkpoint replay is done and before the runner is built,
and returns IndexedTaskInstanceNotStarted as the outcome: no code run, no
checkpoint, no callback, and the next attempt runs it. IndexedTaskOutcomes
counts these apart from the indexed tasks that ran and the kill message names
them. A sync indexed task is still handed to the pool afterwards, whose threads
are as many as the calls in flight, so only that pickup remains. The docs page
and the class docstring list it with the other outcomes that fire no callback.

test_an_item_pulled_before_the_kill_but_not_started_does_not_run fails before
with two of three indexed tasks run; the outcomes unit test covers the count.
@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch from f054261 to 0ce65d5 Compare October 9, 2026 16:34
The static checks fail on D401 for IndexedTaskInstance.context_for: its
summary line named what the method returns instead of saying what it does.
@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch from 0ce65d5 to f482fe4 Compare October 9, 2026 16:51
The message claimed "the rest never pulled" also when every item had been
pulled, and said the input was being resolved for a kill that landed before
the run started, where nothing was. IndexedTaskOutcomes._killed_message now
names the items that ran, those pulled before the kill and never started and
those never pulled, each only when the count is not zero, and a kill before
the input was resolved says so.
test_a_kill_after_every_item_was_pulled_claims_no_remainder fails before.
@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch from f482fe4 to cab90c1 Compare October 9, 2026 17:35
@dabla dabla added the skip newsfragment check Skip the newsfragment PR number check label Oct 9, 2026
@dabla

dabla commented Oct 9, 2026

Copy link
Copy Markdown
Contributor Author

I know this is draft, but please update the PR title to say what it's adding/what it allows, not the specific endpoint it allows.

You're right Ash, will update the title to what this PR's functionally adds.

@dabla dabla changed the title Add a POST /xcoms/{dag_id}/{run_id}/{task_id}/keys execution API endpoint Fetch several XCom values of a task instance in one Execution API request Oct 9, 2026
@dabla dabla changed the title Fetch several XCom values of a task instance in one Execution API request Fetch several XCom values of a Taskinstance in one Execution API request Oct 9, 2026
dabla added 11 commits October 10, 2026 07:51
…oint

XComIterable, the result of an iterated task, stores one XCom per iteration
under distinct keys (return_value_0, return_value_1, ...) of the same task
instance. The existing slice endpoint for a mapped task's XComs ranges over
map_index for a single key, the inverse shape, so iterating or slicing an
XComIterable cost one GET per value.

The new endpoint takes a list of keys and returns their values in that order,
None for a key without an XCom, filtered by map_index (-1 by default), in one
database query. Around it:

* XComKeysRequest body model, and an AddXComKeysEndpoint execution API version
  change so older clients keep their contract.
* has_xcom_access reads the optional key from the request's path parameters,
  so the endpoint needs no separate router or dependency.
* GetXComByKeys supervisor message, XComOperations.get_by_keys client method,
  handle_get_xcom_by_keys handler registered with the task supervisor and the
  DAG processor, an AddGetXComByKeys supervisor schema version entry and the
  regenerated schema snapshot.
* BaseXCom.get_by_keys sends the message and deserializes the values, and
  XComIterable iterates and slices through it: one request instead of one per
  value. A single index stays one XCom read.

Stacked on the task iteration PR, which owns XComIterable; only this commit
belongs to the endpoint. Squashed from the earlier history of this branch and
re-based on that PR, with XComIterable's batched reads moved to its current
home in bases/xcom.py.
@dabla
dabla force-pushed the feature/add-xcoms-keys-route branch from cab90c1 to dcf8439 Compare October 10, 2026 06:27

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

Labels

area:API Airflow's REST/HTTP API area:dag-processor area:task-sdk skip newsfragment check Skip the newsfragment PR number check

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants