namanjain24-sudo opened a new pull request, #73375:
URL: https://github.com/apache/airflow/pull/73375

   The task state store is keyed by the positional `map_index`, and it survives 
a clear on purpose, so a cleared task resumes from its last checkpoint. As 
#72725 shows, that stops being safe for a mapped task once the list it expands 
over changes. Clear the upstream task with its downstream, the upstream returns 
`["z", "a", "b", "c"]` instead of `["a", "b", "c"]`, and map index 0 now 
processes `z` while reading the checkpoint written for `a`. The run still 
succeeds.
   
   The list can only change when the task that produces it runs again. So this 
deletes a mapped task's state store rows, for every map index in that run, when 
a task it expands over is cleared, and leaves them alone otherwise:
   
   | What is cleared | Mapped task's state |
   | --- | --- |
   | The task it expands over, alone or with downstream | deleted, all map 
indices |
   | The whole Dag run | deleted |
   | Only the mapped task | kept, it resumes as documented |
   | Unmapped tasks, and mapped tasks over a literal list | kept |
   
   Tasks inside a mapped task group are handled the same way, using the group's 
expansion input. The check runs in `clear_task_instances`, which both 
`dag.clear()` and the clear endpoints go through. It resolves the Dag on the 
same version the cleared task instances will run on, and logs how many rows it 
deleted.
   
   This differs from #72788, which deletes a task's rows on every clear. That 
also covers this case, but it removes the resume-after-clear behaviour that 
`resumable-tasks.rst` documents for unmapped tasks and for mapped tasks whose 
input is unchanged. The trade-off here is that a mapped task whose upstream is 
cleared starts over even if the upstream happens to return the same list again. 
Comparing the lists would need the old values, which are not stored.
   
   The docs now say this in the "Mapped tasks" section of 
`task-state-store.rst` and in the clearing note in `resumable-tasks.rst`. The 
task state store is not in a release yet (3.3.2 does not have it), so I did not 
add a newsfragment.
   
   Tests:
   
   - New tests in `test_cleartasks.py` cover the cases in the table, a mapped 
task group, a second run of the same Dag that is left alone, and `dag.clear()` 
with downstream. With the new call skipped, every case that expects rows to be 
deleted fails. The "only the mapped task" case passes either way, which is the 
point of that case.
   - An end-to-end check that I did not commit: each map index wrote its item 
through the Execution API task state store endpoint, then the task run was 
cleared, then the new task instances read their checkpoints back. After 
clearing the upstream they got `{0: 404, 1: 404, 2: 404}` with this change, and 
`{0: "a", 1: "b", 2: "c"}` without it, which is the stale state from the issue. 
After clearing only the mapped task they got `{0: "a", 1: "b", 2: "c"}` both 
ways.
   - `test_cleartasks.py` passes on SQLite and on Postgres 16 (47 passed). The 
execution API and public API task state store tests pass, as do the clear tests 
in the public task instances API and `test_dag.py`. prek passes, including mypy 
for airflow-core.
   
   closes: #72725
   
   ---
   
   ##### Was generative AI tooling used to co-author this PR?
   
   - [X] Yes (please specify the tool below)
   
   Generated-by: a Gen-AI coding assistant, following [the 
guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions).
 I reviewed the change and ran the checks above.
   


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