dabla opened a new pull request, #73807: URL: https://github.com/apache/airflow/pull/73807
Stacked on #62922 (Task Iteration): only the last commit, bf1b5a44cc, 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: ```python 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? - [X] Yes (please specify the tool below) Generated-by: Claude Code following [the guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions) 🤖 Generated with [Claude Code](https://claude.com/claude-code) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
