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]

Reply via email to