Skip to content

Add XComIterable.flatten() to read an iterated task's pages as one sequence - #73807

Draft
dabla wants to merge 15 commits into
apache:mainfrom
dabla:feature/xcom-iterable-flatten
Draft

dabla wants to merge 15 commits into
apache:mainfrom
dabla:feature/xcom-iterable-flatten

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, c490ee3, belongs to this PR. It will be rebased once #62922 merges.

related: #62922

What

An iterated task whose iterations each return a page of items, for example one list of rows per API call, pushes one XCom per iteration and returns the XComIterable over them. .flatten() on that output gives a downstream task those pages as one Sequence of items:

pages = fetch_page.iterate(page=range(10))  # each iteration returns a list of rows
load_rows(pages.output.flatten())  # one sequence of rows, page boundaries gone

Nested lists, tuples and sets are expanded recursively; strings, bytes and scalars stay single items.

How

  • FlattenedXComIterable speaks in flattened positions. Its length is the flattened_length that IterableOperator.axcom_push tallies from each value as it is pushed, while the value is in memory, so the view knows its length without reading a page. An XComIterable pushed without the tally is counted by reading its pages once.
  • Reads hold the last page fetched and nothing more: a sequential read (iteration, a slice, an iterated task consuming this as its input) fetches every page exactly once, a read within the held page costs nothing, and a jump backwards restarts from the first page. Memory stays bounded by one page.
  • The same cursor serves __getitem__/__iter__ (blocking reads) and aget/__aiter__ (reads through asend, safe on the task's event loop).

Why a separate PR

This was part of #62922, but task iteration does not need it: nothing in the SDK outside bases/xcom.py referred to it, and it had no documentation there. Splitting it out keeps #62922 about running a task over its input, and gives flatten its own review and its own paragraph in the task-sdk docs, included here.

Tests

test_xcom.py covers expansion of lists, tuples, sets, generators and mixed pages, string/bytes/scalar pass-through, negative indices on both paths, the tallied length round-tripping through serde, one page read per page on a sequential read, the no-tally counting pass, and the one-page memory bound. test_serde.py round-trips the flattened view.


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 commented Sep 27, 2026

Copy link
Copy Markdown
Contributor Author

Interaction with #70223 to handle when both land: FlattenedXComIterable inherits __iter__ from XComIterable. Today that iterates by index through __getitem__, so the flattened view yields flattened items. #70223 changes XComIterable.__iter__ to fetch every raw page with one GetXComByKeys request, so once it is in, FlattenedXComIterable must define its own __iter__ over its flattened positions (self[i] for i in range(len(self))), or iterating it would yield the raw pages. Whichever of the two merges second picks this up.

@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch 15 times, most recently from e1ca970 to bfa2dc4 Compare October 2, 2026 07:52
@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch from bfa2dc4 to 2f26f11 Compare October 5, 2026 17:41
Add Iterable Tasks: `.iterate()` and `.iterate_kwargs()` on operators and
`@task`, the counterpart of `.expand()` that processes every item inside
one task instance instead of creating one task instance per item.

- IterableOperator resolves the input by index as `.expand()` does and
  runs the items on AsyncAwareExecutor: sync operators in a thread pool,
  async operators concurrently on one event loop, up to
  `task_concurrency` at a time.
- Each item's return value is pushed as `return_value_<index>`, and the
  task returns an XComIterable, a lazy read-only Sequence over them that
  a downstream `.expand()` or `.iterate()` consumes. Skipped items are
  left out, and downstream tasks with `all_success` are skipped, as with
  a mapped upstream.
- Per-item progress is checkpointed in the task state store (AIP-103),
  tied to the item's input and the attempt that wrote it, so a retry or a
  clear after a failure resumes the items that already succeeded and a
  clear after success runs them all again. Outlet events are replayed
  from the checkpoint.
- XComs and task state written from an item carry its index, and each
  item runs against its own view of the context.
- Deferral, reschedule-mode sensors, TriggerDagRunOperator and
  downstream skipping from an item are rejected with a clear error.
- Documented in task-sdk/docs/mapped-tasks-vs-iterable-tasks.rst.

Co-Authored-By: Tzu-ping Chung <uranusjr@gmail.com>
Co-Authored-By: Copilot <223556219+Copilot@users.noreply.github.com>
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch from 2f26f11 to 54dc925 Compare October 8, 2026 15:13
Two tests of the iterable operator named the reviewer whose thread led to
them instead of saying what they check. A test's name and docstring must
stand on their own: one now says that a succeeded item's outlet events are
merged into the task's accessor, the other that a failed item's retry
callback waits until no sibling rules the retry out.
dabla added 7 commits October 8, 2026 17:51
The item's view of the context swapped ti, task_instance, task_state_store
and outlet_events but not task, so inside execute and in the item's
callbacks context["task"] was the IterableOperator. Under .expand() it is
the unmapped operator of the task instance: context_update_for_unmapped
sets context["task"] next to ti.task, and the rendering in _create_task
already got that. IndexedTaskRunner.indexed_context now swaps the same key.

Two tests fail before the change: the runner test for the indexed context
asserts the item's operator under "task" and the parent's own left alone,
and an operator test checks that execute and on_success_callback of every
item see the unmapped operator there, the same object as ti.task.
The docs page presented iteration as a way to keep large results out of
the metadata database through a custom XCom backend. Each item's result is
also written to its checkpoint, which goes to the task_state_store table
for the store's retention unless a state_store_backend is configured, so
with only a custom XCom backend the payloads move from one table to the
other. The triggerer bullet, the "XCom backend" row of the comparison
table and the class docstring now say so and point at state_store_backend.
A task instance sends one event per asset, so items of an iterated task
that emit to the same asset are merged into that one event: extra keeps
what the last item to finish wrote, partition keys and alias events
accumulate. .expand() sends one event per mapped task instance, and a
consumer reading triggering_asset_events sees all of them. The comparison
table gets an "Asset events" row and _merge_outlet_events a note, since a
per-file Metadata pattern ported from .expand() produces one event here.
.expand() never sees a string or another scalar from upstream: the
upstream's _push_xcom_if_needed raises UnmappableXComTypePushed when it
has a mapped dependant. An IterableOperator is not a MappedOperator, so
iter_mapped_dependants does not find it and that check never fires; an
upstream returning a JSON string made .iterate() succeed over one wrong
item. Source.from_argument now applies the same rule to the value it
resolved from an XComArg, with the is_mappable_value the push check uses,
and XComForMappingNotPushed for None, as the push would have raised.

Tests cover strings, bytes, scalars and None on Source, and .iterate()
over an upstream returning a JSON string, which fails before any item runs.
With K failed items and a retry_policy, the policy was evaluated K times
in _failure_for_the_runner to choose the exception for the runner, once
more in _task_will_retry to pick the failed items' callbacks, and once
more by the runner itself. A policy that calls a model, such as
common.ai's LLMRetryPolicy, pays for each call and may answer differently
each time, so the callbacks could announce a retry the runner then did not
take, which deferring them was meant to prevent.

_failure_for_the_runner now keeps the decision taken for the exception it
chose, and _task_will_retry reuses it; the policy is only evaluated there
for an exception that was never chosen (a single failure). Two tests fail
before: a counting policy is evaluated once per failed item, and a policy
that alternates its answers gives the callbacks the same outcome as the
exception handed over.
…ff the loop

Only run() went to the thread pool for a sync item: IndexedTaskRunner's
enter and exit, and so the item's on_success_callback and
on_skipped_callback fired from the exit, ran on the event loop thread. A
notifier reading a connection or a variable there while a sibling was
inside an async SDK call got DeadlockImminentError, a BaseException that
_run_task_state_change_callbacks does not catch: the item had done its
work, was checkpointed UP_FOR_RETRY anyway, and the task failed without a
retry with a message blaming an async sub-task. The user wrote an ordinary
sync operator; the operator put it on a thread.

A sync item now enters and exits its runner inside the function handed to
the thread pool, so its callbacks run where a sync SDK call waits for the
comms lock; async items stay on the loop. A sync item whose coroutine was
cancelled while its thread went on reports nothing from its exit
(IndexedTaskRunner.cancel), since it gets no checkpoint and runs again on
the next attempt. on_kill catches BaseException per sub-operator, so one
failure no longer skips the rest, and runs the kills off the loop thread:
in a thread of their own when the SIGTERM handler calls it while the loop
runs, through asyncio.to_thread from _run_tasks. The failure message for a
DeadlockImminentError says which thread the call was made on, and the
class docstring says where each callback runs.

Tests pin the thread of execute and of every callback for a sync and an
async item, the kill reaching the second sub-operator when the first
raises DeadlockImminentError, the kill running off a running loop, and a
cancelled sync item firing no callback when its thread finishes.
On SIGTERM the runner calls on_kill() once and the items in flight are
killed, but nothing stopped the iteration from pulling more: each killed
item came back as a failure and freed its slot, _fill_pending submitted
the next item into it, and new items with their remote jobs kept starting
until the supervisor escalated to SIGKILL. Those never got an on_kill.

IterableOperator keeps a stop flag that on_kill() sets first.
AsyncAwareExecutor.imap_unordered takes it as ``stop`` and asks it before
and after every pull: once set nothing more is submitted, an item pulled
at that moment included, and what was already submitted drains. After the
loop, before the other outcomes, _run_tasks raises AirflowTaskTerminated
saying how many items ran and that the rest never started: killed items
returning normally would otherwise have written the completion marker and
returned an XComIterable over XComs that do not exist. The task fails
without a retry, as the runner treats a terminated task, and no marker is
written, so a later clear resumes from the checkpoints of the items that
did finish.

Tests: the executor stops pulling once ``stop`` says so and still drains
what it submitted; a kill from inside the first of four items runs no
other item, leaves no marker and only that item's checkpoint; a kill
while the input is resolved runs no item and is not an empty-input skip.
@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch from 54dc925 to 772416a Compare October 8, 2026 17:24
dabla added 2 commits October 8, 2026 19:47
IterableOperator carried seven attributes for what one run remembers
while it is going: the sub-operators in flight with their lock and the
ones already killed, the stop flag, the runners of the failed items, the
resolved input and the retry policy's decision for the exception handed to
the runner, plus a __deepcopy__ override to give every copy fresh ones.

They move into IterationState, in the same module. A copy of the state is
a fresh one, so the override goes, and the operator starts every run with
a fresh state too. The state exposes what the operator means rather than
its containers: register/unregister and take_in_flight for the kill,
request_stop/stop_requested for the executor's stop, note_failed and the
failed_runners tuple, keep_decision/decision_for, and length for the
resolved input; the operator makes no reach-through call into it.
IndexedTaskRunner registers through a SubOperatorRegister protocol instead
of a dict and a lock.

Tests cover the state on its own: registration until unregistered, keying
by identity for operators that compare equal, take_in_flight handing each
operator out once, the stop flag, the order of failed runners, a kept
decision answering for that exception object only, length unknown until
resolved, and a deep copy being fresh while the original keeps its state.
BaseOperator.__init__ runs under _apply_defaults, which fills every
parameter of its signature the call leaves out from the DAG's default_args.
IterableOperator forwarded the retry, scheduling and pool settings of the
wrapped operator but left the five on_*_callback parameters and
pre_execute/post_execute out, so a DAG with
default_args={"on_failure_callback": notify} gave the iterated task that
callback too, while the items got it through the wrapped operator's partial
kwargs: with three failed items notify ran four times, the last one with
the parent's context, where .expand() and .partial(on_failure_callback=
notify) run it three times. The class docstring and the docs page say the
iterated task carries no callbacks of its own.

They are now passed as None explicitly, so _apply_defaults leaves them
alone and both paths agree with the docs. The test builds a DAG with all
seven in default_args and checks the iterated task has none while an
unmapped item has them; it fails before the change.
@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch from 772416a to 980397c Compare October 8, 2026 18:15
on_kill() took a thread only when the event loop was running on the calling
thread. The runner's SIGTERM handler can also arrive while the loop is
paused between two run_until_complete calls, as imap_unordered hands a
result to the consumer; the kills then ran inline on the main thread, and a
synchronous SDK call in a sub-operator's on_kill waited for the comms lock
a parked asend held, which only the paused loop could release: the loop
never resumed, the remaining on_kills never ran, and SIGKILL ended it.

on_kill() now always sets the stop flag, takes the in-flight snapshot and
kills in a daemon thread of its own, where the call waits its turn in both
cases while the loop goes on. _run_tasks no longer goes through on_kill():
it kills through asyncio.to_thread as before, awaited while the loop runs.

A new test calls on_kill() with no loop running on the main thread, as in
the paused window, and checks the kill ran off it; the tests that asserted
the kill right after the call wait for the thread.
@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch from 980397c to 8d9cf89 Compare October 8, 2026 21:53
…nd the async item's own timeout too

Two on_kill call sites still ran on the loop thread. IndexedTaskRunner.in_flight
killed the operator the parent's execution timeout struck as the timeout unwound,
and _execute_async_task killed an async item that ran out of its own
execution_timeout; both on the loop thread, where an on_kill that cancels a
remote job through a sync hook raises DeadlockImminentError. The first site let
it through its `except Exception`, the second was unguarded, so the parent's
timeout became an AirflowFailException: no retry, and the job left running.

The operator the parent's timeout strikes now stays registered as the timeout
unwinds (the timeout only lands on async items: the signal reaches the main
thread, where the loop runs, while sync items sit in worker threads), and
IterableOperator._run_tasks kills it off the loop with the others through
asyncio.to_thread. An async item's own timeout kills it through
`await to_thread(task.on_kill)`.

Tests fail before: the runner leaves the struck operator registered and does
not kill it in place; the async item's own timeout kills it off the loop thread;
an async item raising the parent's timeout inside its coroutine, with an
on_kill that raises on the loop thread, still ends in AirflowTaskTimeout with
both items killed.
@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch from 8d9cf89 to 7366809 Compare October 8, 2026 22:00
dabla and others added 2 commits October 9, 2026 00:03
…ync item

Every sync item with an execution_timeout went through _run_execute_callable
in its worker thread, which sent SetExecutionTimeout to the supervisor again
and tried TimeoutPosix, which only works on the main thread: the supervisor's
hard-kill deadline drifted to "timeout after the last item started", up to a
full timeout late, and the TimeoutPosix warning was logged once per item.
With 4 sync items and a 30 s timeout: 4 sends, 4 warnings.

_run_execute_callable takes enforce_timeout, False for an IndexedTaskInstance:
the items run in worker threads under the parent's own limit, which the
parent enforces and reported once. The async items already handled their
limit themselves. A test runs three sync items with a 30 s timeout through
the parent's _run_execute_callable and checks a single SetExecutionTimeout;
it sees four before the change.
…quence

An iterated task whose iterations each return a page of items (a list of rows
per API call) pushes one XCom per iteration and returns the XComIterable over
them. flatten() gives the downstream task those pages as one Sequence of items:
nested lists, tuples and sets are expanded recursively, strings, bytes and
scalars stay single items.

FlattenedXComIterable speaks in flattened positions. Its length is the
flattened_length that IterableOperator.axcom_push tallies from each value as
it is pushed, so the view knows it without reading a page; an XComIterable
pushed without the tally is counted by reading its pages once. Reads hold the
last page fetched and nothing more: a sequential read fetches every page
exactly once, a read within the held page costs nothing, and a jump backwards
restarts from the first page, so memory stays bounded by one page.

Split out of the task iteration PR, which does not need it, so it can be
reviewed and documented on its own.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@dabla
dabla force-pushed the feature/xcom-iterable-flatten branch from 7366809 to 6ee29e6 Compare October 8, 2026 22:04

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