ColtenOuO opened a new pull request, #72007: URL: https://github.com/apache/airflow/pull/72007
### Sumarry Fix the Execution API's XCom sequence-read endpoints so a mapped task/task-group instance that hasn't pushed an XCom yet (still `up_for_reschedule`, or skipped) doesn't shift the values of every mapped instance after it. closes: #40321 ### The bug `get_mapped_xcom_by_index`, `get_mapped_xcom_by_slice`, and `head_xcom` (in `airflow-core/src/airflow/api_fastapi/execution_api/routes/xcoms.py`) all built their result by querying XCom rows ordered by `map_index` and then using the *position in that result set* as if it were the `map_index` itself (via SQL `OFFSET`/`LIMIT`/`slice()`, with `map_index` dropped from the response). That's correct only when every mapped instance has already pushed an XCom. If one hasn't — e.g. a sensor inside a dynamically mapped `@task_group` is still `up_for_reschedule` — there's no row for its `map_index`, so the query silently returns one fewer row and every later `map_index`'s value shifts down by one position. A downstream task reading `mapped_group.output[i]` (or iterating over it) can get another instance's value instead of the one it asked for, or hit an `IndexError` reading past the (now short) end. This is the same root cause reported in #40321 in 2024 against the Airflow 2.x `LazyXComSelectSequence` implementation. The AIP-72 rewrite reimplemented the same "position == map_index" assumption in the new task-sdk/execution-API code path instead of fixing it, so it's still present on `main` today. ### Confirming the bug is still there, and that the fix resolves it Before writing the fix, I reproduced the bug against current `main` with a test scenario matching the issue: a task mapped over 4 values, where `map_index=1` never gets an XCom row (simulating a mapped instance that's still `up_for_reschedule` and hasn't pushed anything). Querying the slice endpoint for all 4 values returned only 3 — `["f", "o", "b"]` — with the value that actually belongs to `map_index=3` silently shifted into `map_index=1`'s position. The `HEAD` endpoint's count also undercounted (3 instead of 4). This matches exactly what the issue reports: values read by a downstream task line up with the wrong `map_index`, and the reported length is wrong too. Two of the existing tests in `test_xcoms.py` (`test_xcom_get_with_offset`, `test_xcom_get_with_slice`) already built this exact gapped scenario, but asserted the *compacted* 3-item result as the expected outcome — meaning the buggy behavior was pinned as correct by the test suite rather than being an untested gap. After the fix, I reran the same scenario: the slice endpoint now returns the full 4-item `["f", None, "o", "b"]`, with `map_index=3`'s value correctly in position 3, and `HEAD`'s count correctly reports 4. I updated the two pinning tests to assert this corrected behavior and added `test_xcom_get_with_slice_and_count_unfinished_mapped_task` as an explicit regression test for the issue's scenario (a task instance that never pushed an XCom at all, as opposed to one whose pushed value happened to be `None` — those are different situations that both need to resolve correctly). ### The fix Added `_get_mapped_length()`, which resolves the *logical* number of mapped instances for `(dag_id, run_id, task_id)` from the `TaskInstance` table (every mapped instance gets a row at expansion time, whether or not it has run yet — see `TaskMap.expand_mapped_task`), independent of how many XCom rows exist. When a task is mapped in this run, the three endpoints now key results by the *actual* `map_index` instead of result position, filling in `None` for any `map_index` that has no XCom row yet: - `get_mapped_xcom_by_index`: `offset` resolves against the logical `0..mapped_length-1` range; a valid-but-missing index returns `None` (200), not a 404 — only an out-of-range index 404s. - `get_mapped_xcom_by_slice`: builds a `map_index -> value` dict, fills gaps with `None`, then applies the slice with plain Python slicing. - `head_xcom`: reports the logical mapped length in `Content-Range`, not the XCom row count. Two cases are deliberately left on the old (pre-fix) behavior, since there's no well-defined "logical length" for them: - `include_prior_dates=True` on the slice endpoint can span multiple dag runs, each with its own (potentially different) mapped length. - An unmapped task (`map_index=-1`, no expansion) has no mapped-length concept at all — it falls through to the original row-position logic unchanged. `get_xcom`'s `offset` query-parameter branch (a separate, older code path) is intentionally untouched — no task-sdk client currently calls it with `offset` set, and it has its own pre-existing "skip None values" semantics that aren't part of this bug. ### Blast radius: what else calls these endpoints, and what changes for callers I checked who actually consumes these three routes and what the response-shape change means for them, since this touches a fairly central part of the Execution API. **Consumers.** `task-sdk`'s `LazyXComSequence` (`task-sdk/src/airflow/sdk/execution_time/lazy_sequence.py`) is the only real consumer, reached through `XComOperations` in `task-sdk/src/airflow/sdk/api/client.py`. I grepped the rest of the repo (`providers/`, the UI, everywhere else) for these route paths and the relevant symbol names and found no other caller — nothing outside task-sdk talks to these three endpoints directly. `client.py` itself is a thin pass-through and doesn't assume anything about the old semantics, so it needed no changes. **User-visible behavior change.** This is the one worth calling out explicitly: a mapped task's `.output` used to silently compact away gaps when iterated, sliced, or expanded over downstream (`for v in t.output`, `.expand(x=t.output)`, etc.) — a missing map_index just meant the list was shorter than the actual mapped count, with everything after it shifted into the wrong position. After this fix, `.output` is the correct length and gaps show up explicitly as `None` at their true position. This is the intended correctness fix, but it's a behavior change any DAG author relying (knowingly or not) on the old compaction would notice — e.g. `.expand(x=t.output)` now creates an instance with `x=None` for the gap instead of one fewer instance. Flagging this for review since it's a semantic change, not just a bugfix that's invisible from the outside. **Performance.** The original version of this fix always fetched every XCom row for the task into a Python dict once a task was detected as mapped, even when there was no gap at all — regressing the existing SQL-side `OFFSET`/`LIMIT`/`slice()` pagination for the common case. Since XCom values are stored in an unbounded JSON column (only the *count* of mapped instances is capped, by `core.max_map_length`, default 1024), that could mean loading a large number of large payloads into API-server memory for a single-item slice request. Fixed by adding one cheap `COUNT` query to detect whether a gap actually exists before falling into the full-fetch path — the common case (no gaps) keeps the original SQL pagination; only a genuinely sparse sequence pays the cost of loading every row, which is unavoidable since placing values at their true `map_index` requires knowing all of them. Added `test_xcom_get_with_slice_mapped_task_without_gap_uses_sql_pagination` to cover the no-gap path exp licitly, since the existing gap-scenario tests wouldn't have caught a regression here. ### Testing - Updated `test_xcom_get_with_offset` / `test_xcom_get_with_slice` to assert the corrected sparse result instead of the compacted one, and extended their index/slice boundaries to cover the full logical range. - Added `test_xcom_get_with_slice_and_count_unfinished_mapped_task`: the regression test described above, covering both the slice and count (`HEAD`) endpoints. - Added `test_xcom_get_with_slice_mapped_task_without_gap_uses_sql_pagination`, covering the no-gap performance path described above. - Added two tests confirming unmapped tasks keep the pre-fix (row-position) behavior unchanged. - Added `test_getitem_index_sparse_gap_returns_none` in task-sdk's `test_lazy_sequence.py`, confirming the client resolves a gap to `None` rather than raising `IndexError`. Traced this through `XCom.deserialize_value` / `airflow.sdk.serde.deserialize`, which already short-circuits on `None`, so no task-sdk source change was needed. - Confirmed the existing `versions/v2025_04_28` API-version contract test (a different code path, untouched by this fix) still passes unchanged. - `ruff` and `mypy` (via `breeze run mypy`) pass on all changed files. --- ##### Was generative AI tooling used to co-author this PR? - [X] Yes — Claude Code (Sonnet 5) -- 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]
