dabla opened a new pull request, #73790:
URL: https://github.com/apache/airflow/pull/73790

   Stacked on #62922 (Task Iteration): only the last commit, 1c508ce342, 
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 the `XCOM_SEQUENCE_CHUNK_SIZE` 
(32) 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 module constant for now, 32 items, which bounds memory at 32 values of 
whatever size the XComs are. If reviewers prefer it configurable it can become 
a `[core]` option in a follow-up commit; I kept the config surface out of this 
PR until that is decided.
   
   ## Tests
   
   `test_lazy_sequence.py` now asserts slice requests: 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?
   
   - [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