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]
