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]
