Skip to content

Read a mapped upstream's XCom sequence in chunks instead of one item per request - #73790

Draft
dabla wants to merge 43 commits into
apache:mainfrom
dabla:feature/xcom-sequence-chunked-reads
Draft

dabla wants to merge 43 commits into
apache:mainfrom
dabla:feature/xcom-sequence-chunked-reads

Conversation

@dabla

@dabla dabla commented Sep 27, 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): only the last commit, b0859ce, belongs to this PR. It will be rebased once #62922 merges.

related: #62922

What

LazyXComSequence, the value a mapped upstream's XComArg resolves to, fetched one item per GetXComSequenceItem request. Iterating it, or an iterated task reading it by index (see #62922, where .iterate() resolves its input by index), cost one supervisor round trip to the API server per item, which is the dominant cost for large inputs.

Reads now come in chunks:

  • An index outside the slice held fetches chunk_size items starting there with one GetXComSequenceSlice request, the message the slice path already used, and the consecutive reads that follow are served from it. Sequential access costs one request per chunk instead of one per item, on both the sync path (__getitem__/__iter__ through send) and the async one (aget/__aiter__ through asend), and at most one chunk is held in memory. A jump backwards fetches again from there.
  • A negative index no longer needs its own per-item request: it is mapped from the end through the cached length, so one GetXComCount at most.
  • The client no longer sends GetXComSequenceItem at all. The request handler and the API endpoint keep serving it, so nothing else changes.

The sync and async paths share the request builder and the response parser, so they cannot drift.

Why

For a 17,000-item mapped upstream consumed by .iterate(), this is about 530 requests instead of 17,000. It is the "XCom slice optimisation" left out of #62922 on purpose, so that PR stays a pure simplification with unchanged behaviour.

Chunk size

A new option, [core] xcom_sequence_chunk_size, default 32, described in config.yml next to max_map_length and mentioned in the iteration docs. It bounds memory at that many values of whatever size the XComs are, and larger values mean fewer round trips. LazyXComSequence reads it as its default and also accepts an explicit chunk_size.

Tests

test_lazy_sequence.py now asserts slice requests, with the chunk size taken from the config option and given explicitly: one per chunk while iterating, none for a read within the held chunk, a new request on a jump backwards, and negative indices through the cached count on both paths. The mapped-operator fake supervisor honours slice bounds, which it previously ignored.


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

@dabla dabla added this to the Airflow 3.4.0 milestone Sep 27, 2026
@dabla dabla self-assigned this Sep 27, 2026
@dabla dabla mentioned this pull request Sep 27, 2026
1 task done
@dabla
dabla force-pushed the feature/xcom-sequence-chunked-reads branch from 1c508ce to 5302e88 Compare September 27, 2026 12:41
@dabla
dabla force-pushed the feature/xcom-sequence-chunked-reads branch from 5302e88 to 89b6cc1 Compare September 27, 2026 17:32
@dabla

dabla commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

Restacked on #62922 after XComIterable.flatten() moved out to its own PR; still one commit, 89b6cc1.

@dabla
dabla force-pushed the feature/xcom-sequence-chunked-reads branch 17 times, most recently from e7cfaef to 50d3d03 Compare October 2, 2026 07:52
@dabla
dabla force-pushed the feature/xcom-sequence-chunked-reads branch 2 times, most recently from ba13cca to ea231c5 Compare October 8, 2026 15:13
@dabla
dabla force-pushed the feature/xcom-sequence-chunked-reads branch from 9df820e to 8115f75 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/xcom-sequence-chunked-reads branch from 8115f75 to 739ffb6 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/xcom-sequence-chunked-reads branch from 739ffb6 to 8bcec79 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/xcom-sequence-chunked-reads branch from 8bcec79 to 584fc0e 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/xcom-sequence-chunked-reads branch from 584fc0e to aa79f39 Compare October 9, 2026 17:35
@dabla dabla added the skip newsfragment check Skip the newsfragment PR number check label Oct 9, 2026
dabla added 11 commits October 10, 2026 07:51
…per request

LazyXComSequence fetched one item per GetXComSequenceItem request: iterating a
mapped upstream's results, or an iterated task reading them by index, cost one
supervisor round trip per item, which is the dominant cost for large inputs.

Reads now come in chunks. An index outside the slice held fetches chunk_size
items starting there with one GetXComSequenceSlice request, the message the
slice path already used, and consecutive reads are served from it. Sequential
access costs one request per chunk instead of one per item, both through
__getitem__/__iter__ (send) and aget/__aiter__ (asend), and at most one chunk
is held. A jump backwards fetches again from there.

The chunk size is the new [core] xcom_sequence_chunk_size option, default 32,
documented in config.yml and in the iteration docs; a LazyXComSequence can
also be given one explicitly.

A negative index no longer needs its own per-item request: it is mapped from the
end through the cached length, one GetXComCount at most. The client no longer
sends GetXComSequenceItem at all; the request handler keeps serving it.
@dabla
dabla force-pushed the feature/xcom-sequence-chunked-reads branch from aa79f39 to 15c4964 Compare October 10, 2026 06:25

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.

1 participant