dabla commented on PR #62922:
URL: https://github.com/apache/airflow/pull/62922#issuecomment-5855751286

   ### Rework: `.iterate()` now resolves its input by index, the way 
`.expand()` does
   
   Commits 79e7a3364a and 9a75b7898a. The goal is to make this PR easier to 
review, to simplify the code, to avoid unnecessary duplication, and to reuse 
what already exists as much as possible.
   
   **Why.** `IterableOperator` walked its input through 
`iter_values()`/`aiter_values()`, a streaming cross product that existed 
nowhere else in Airflow and had to be kept in a sync and an async flavour so 
the two could not drift. The length was only known after the stream had been 
drained, recorded as a side effect by `count()`/`async_count()`, and the sync 
flavour was never called at run time. Every input `.iterate()` accepts is a 
sized sequence once pulled, so it can do what `.expand()` already does: know 
the length up front and pick each item by index.
   
   **What changed.**
   - `ExpandInput.aresolve(context)` pulls every source once and returns 
`Resolved(length, aget)`. `DictOfLists` is the cross product, `ListOfDicts` the 
mapping at that position, `Decorated` wraps it in `op_kwargs`.
   - The cross-product arithmetic is extracted into `index_for_each_field()`, 
now shared by `_expand_mapped_field()` (`.expand()`) and `aresolve()` 
(`.iterate()`) and tested on its own, so both hand a sub-task the same item for 
a position.
   - `Source` reads one argument by index without blocking the loop thread: a 
value's own `alen()`/`aget()` when it has them (`XComIterable`, 
`LazyXComSequence`, which gains an `alen()` twin of `__len__`), else a worker 
thread.
   - `XComIterable` carries a `flattened_length` that 
`IterableOperator.axcom_push` tallies from each value as it is pushed, so 
`FlattenedXComIterable` knows its length without reading a page. It is now a 
cursor over the last page fetched: a sequential read fetches every page exactly 
once, where indexing the old class re-walked the pages from the start on every 
call.
   - Removed from the PR's API surface: `XComArg.iter_values`/`aiter_values`, 
`ResolveMixin.iter_values`, the `ExpandInput` sync/async twins, `count()`, 
`async_count()`, `aiterate()` and `_to_iterable()`.
   
   Net effect: 150 fewer lines of source.
   
   **Behaviour change.** `iterate_kwargs([xcom1, xcom2])` now follows 
`expand_kwargs`, one mapping per XComArg, where it used to flatten each into 
several.
   
   **Supervisor calls.** Unchanged for in-memory and plain XCom inputs. A 
mapped upstream costs one `GetXComCount` per source, since the length is needed 
up front; its per-item reads are unchanged. A cross product no longer 
re-resolves a later source once per item of an earlier one.
   
   **Note on the static checks job.** The failing hook regenerates the Java SDK 
dependency verification metadata; this PR does not touch the Java SDK. The 
branch is behind `main`, whose newer Java build the metadata file already 
reflects, so the hook prunes entries the branch's build does not reference. It 
clears on the next rebase onto `main`.
   


-- 
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