vasilisnasopoulos opened a new issue, #72725:
URL: https://github.com/apache/airflow/issues/72725

   ## Summary
   
   `task_state_store` rows are keyed by `map_index` — a **positional** index 
into the expanded upstream list — and they survive a clear-and-rerun by design.
   
   If the upstream list changes between the first run and the clear-and-rerun 
(an insertion, a deletion, a reordering — routine for "list the new files since 
yesterday"), map index *N* on the second attempt is a **different item** than 
map index *N* on the first, and it silently inherits the earlier item's state.
   
   The run reports success.
   
   ## Where
   
   - `airflow-core/src/airflow/models/task_state_store.py` — unique on 
`(dag_run_id, task_id, map_index, key)`
   - `airflow-core/src/airflow/state/metastore.py` — `_get` / `_store` / 
`_delete` match on `(dag_id, run_id, task_id, map_index, key)`; the async twins 
do the same
   - `task-sdk/src/airflow/sdk/execution_time/task_runner.py` — builds 
`TaskScope(dag_id, run_id, task_id, map_index)`
   
   Airflow's own clear path is careful about identity elsewhere: 
`clear_task_instances` regenerates the TaskInstance uuid (`taskinstance.py`, 
`self.id = uuid7()`), and a rerun task's XComs are deleted before re-execution 
(`task_runner.py`). Only the state store is addressed positionally.
   
   ## Reproduction
   
   A mapped task that checkpoints per item, run once over `["a","b","c"]`, then 
**cleared and rerun in place** (what the UI "Clear task" button does, via 
`clear_task_instances`) with the upstream now returning `["z","a","b","c"]`:
   
   ```
   --- ATTEMPT 1 (items=[a,b,c]) ---
   OBSERVE map_index=0 item='a' checkpoint_read=None -> wrote 'a'
   OBSERVE map_index=1 item='b' checkpoint_read=None -> wrote 'b'
   OBSERVE map_index=2 item='c' checkpoint_read=None -> wrote 'c'
   
   --- ATTEMPT 2 (clear-and-rerun in place, items=[z,a,b,c]) ---
   OBSERVE map_index=0 item='z' checkpoint_read='a'  [MISPAIRED]
   OBSERVE map_index=1 item='a' checkpoint_read='b'  [MISPAIRED]
   OBSERVE map_index=2 item='b' checkpoint_read='c'  [MISPAIRED]
   OBSERVE map_index=3 item='c' checkpoint_read=None -> wrote 'c'
   [rerun] final dagrun state: success
   ```
   
   Every pre-existing checkpoint was handed to the wrong item. In the 
documented external-job-resumption pattern this means the instance processing 
`z` reattaches to the external job submitted for `a`.
   
   The DAG:
   
   ```python
   @dag(schedule=None, start_date=LOGICAL_DATE, catchup=False)
   def d():
       @task
       def make_items():
           return _items()          # reads a file; changed between the two 
attempts
   
       @task
       def process(item, **context):
           store = context["task_state_store"]
           ck = store.get("checkpoint")
           if ck is None:
               store.set("checkpoint", item)
           else:
               print(f"map_index={context['ti'].map_index} item={item!r} 
checkpoint_read={ck!r}")
   
       process.expand(item=make_items())
   ```
   
   ## Negative controls
   
   **A — same mechanism, list unchanged** (the documented survive-a-clear case):
   
   ```
   OBSERVE map_index=0 item='a' checkpoint_read='a'  [OK]
   OBSERVE map_index=1 item='b' checkpoint_read='b'  [OK]
   OBSERVE map_index=2 item='c' checkpoint_read='c'  [OK]
   ```
   
   **B — shifted list but a different DagRun** (the near miss that correctly 
does not trigger — `run_id` in the key does isolate):
   
   ```
   OBSERVE map_index=0 item='z' checkpoint_read=None -> wrote 'z'
   OBSERVE map_index=1 item='a' checkpoint_read=None -> wrote 'a'
   ```
   
   ## The silence
   
   Driver exit code 0. All task instances `success`. DagRun `success`. Grepping 
the full run log for `[warning`, `[error` and `Traceback` yields nothing 
related. Nothing notices that the expansion width or contents changed, and 
nothing notices that a checkpoint written under a different input is being 
served.
   
   ## What the docs say
   
   `airflow-core/docs/core-concepts/task-state-store.rst`:
   
   > When a task is dynamically mapped (`task.expand(...)`), each map index has 
its own task state store namespace.
   
   presented as an isolation guarantee, with no caveat that the index is 
positional or that it can name a different item after re-expansion. The same 
page states that surviving an operator-initiated clear is the intended use, so 
the persistence is not the problem — what persists is addressed by position.
   
   Separately, `docs/authoring-and-scheduling/dynamic-task-mapping.rst` states 
that **the order of expansion is not guaranteed** and that mapped tasks are 
assigned an integer index. The two statements are never joined.
   
   The class comment in `models/task_state_store.py` reasons about retries 
("retries of the same task share the same rows — that is the point") and about 
different DAG runs having different `dag_run_id`. The re-expansion case is not 
considered.
   
   The sibling `AssetScope` docstring **does** carry the equivalent warning, 
that `name` and `uri` are not guaranteed to be unique over time. `TaskScope` 
has none.
   
   The documented custom-backend example builds the external object path as 
`f"airflow/task-store/{scope.dag_id}/{scope.run_id}/{scope.task_id}/{scope.map_index}/{key}"`,
 so the positional identity propagates into user object stores and the 
mispairing follows it out of the metadata database.
   
   ## Note on reproducing this
   
   `dag.test()` cannot show it: `get_or_create_dagrun` **deletes** the existing 
DagRun, and the `ON DELETE CASCADE` foreign key wipes `task_state_store` with 
it, so every read comes back `None`. That looks like the store is safe and is 
not. The reproduction above uses `clear_task_instances` — the same function the 
public clear endpoint calls — and keeps the DagRun row.
   
   ## Not tested
   
   - No scheduler/executor/API-server deployment. Tasks ran through Airflow's 
own `_run_task` in-process against SQLite; the Postgres/MySQL upsert dialects 
in `metastore.py` were not exercised.
   - The re-run loop in the harness is mine, modelled on the body of 
`dag.test()`, because no public helper re-runs an existing DagRun in process. 
The **clear** and the **task execution** are Airflow's code.
   - Custom (non-metastore) state backends and the external-ref serialization 
path.
   - The async `a*` variants — identical code shape, presumed identical, 
unexercised.
   - `map_index_template`-named indices, which change the UI label but not the 
stored `map_index`.
   - The list-shrink direction: rows at now-nonexistent higher indices persist 
until expiry and would be picked up if the list grows again.
   - `clear_on_success=True`, which would mask this for successful tasks but 
not for tasks cleared from a failed state.
   - Backfill, and multi-argument `expand` / `expand_kwargs` / `.concat()`.
   
   I have not checked whether this has been raised before; a pointer to an 
existing issue is welcome and I am happy to close this in favour of one.
   
   ## Version
   
   Airflow `3.4.0`, source at `33d9915992980e1a082f665d33861f39171393da`, 
installed from source, Python 3.12, SQLite metastore.
   


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