dabla commented on code in PR #62922:
URL: https://github.com/apache/airflow/pull/62922#discussion_r4093217191
##########
task-sdk/src/airflow/sdk/bases/xcom.py:
##########
@@ -564,3 +565,294 @@ def delete(
map_index=map_index,
),
)
+
+
+class XComIterable(Sequence):
+ """
+ An iterable that lazily fetches XCom values one by one instead of loading
all at once.
+
+ The class has two sides. The *producing* task builds it and grows it with
:meth:`append` /
+ :meth:`aappend`, each call pushing one more ``return_value_<index>`` XCom,
before returning it as
+ the task's result. Everything *downstream* (``.iterate()``, ``.expand()``,
a plain ``xcom_pull``)
+ only ever reads it, which is why the class implements the read-only
+ :class:`collections.abc.Sequence` rather than ``MutableSequence``: once
handed over it is a fixed
+ view of the values already pushed, and the two append methods are not part
of that contract.
+
+ Negative indices are not supported: every element is a remote fetch, and
resolving a negative
+ index against a lazily counted stream would cost a full walk just to find
the end.
+ """
+
+ def __init__(
+ self,
+ task_id: str,
+ dag_id: str,
+ run_id: str,
+ map_index: int | None = None,
+ length: int | None = None,
+ ):
+ self.task_id = task_id
+ self.dag_id = dag_id
+ self.run_id = run_id
+ self.map_index = map_index
+ self.length = length or 0
+ self._index = self.length
+
+ def __iter__(self) -> Iterator[Any]:
+ return _XComIterator(self)
+
+ def __len__(self) -> int:
+ return self.length
+
+ def __getitem__(self, key: int | slice) -> Any | Sequence[Any]:
+ """Allow direct indexing so this works like a sequence."""
+ from airflow.sdk.execution_time.xcom import XCom
+
+ if isinstance(key, slice):
+ # TODO: This issues one XCom.get_one call per element — N
round-trips for a full slice.
+ # XComIterable stores results under distinct keys (return_value_0,
return_value_1, …)
+ # with the same map_index, so the existing GetXComSequenceSlice
endpoint (which ranges
+ # over map_index for a single key) cannot be reused. A new POST
endpoint that accepts
+ # a list of keys and returns values in a single query is needed;
once that lands, replace
+ # this loop with a single batched fetch.
+ start, stop, step = key.indices(len(self))
+ return [self[i] for i in range(start, stop, step)]
+
+ if not (0 <= key < self.length):
Review Comment:
Fixed in 94219a4721. Both classes now honour negative indices instead of
rejecting them, so `result[-1]` is the last element on either view; the
flattened one already walked the stream for its bounds check, so resolving from
the end is free, and its `aget` discovers the count on the loop.
`append`/`aappend` are removed from `XComIterable`, so it is now a genuinely
read-only `Sequence`: the operator pushes per-index values through
`axcom_push`, and nothing in this PR needs a producer-side API. A dedicated
producer-side class (one that mutates, is deliberately not a `Sequence`, and
deserializes as a plain `XComIterable`) is useful for paginated sources, but it
is out of scope here and will follow in a separate PR once this one is merged,
the same way deferred-operator support inside `IterableOperator` was split out.
The paginated example in the PR description that relied on `append()` is
removed for the same reason.
---
Drafted-by: Claude Fable 5.1; reviewed by @dabla before posting
--
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]