Repository navigation
Reset a mapped task's state store when its expansion is recomputed - #73375
namanjain24-sudo wants to merge 1 commit into
Conversation
The task state store is keyed by the positional map_index, and it deliberately survives a clear so a cleared task resumes from its last checkpoint. For a mapped task that is only safe while the expanded list stays the same. When a task the mapped task expands over is cleared, it runs again and can return a different list, and each map index then read the state written for whichever item used to sit at that position. The run still succeeded, so nothing surfaced the mix-up. clear_task_instances now deletes the state store rows of every mapped task, and every task in a mapped task group, that expands over one of the cleared tasks, for all map indices of that run. Clearing only the mapped task keeps its rows, since its input does not change, so the documented resume-after-clear behaviour is unchanged there. The rows are looked up on the same Dag version the cleared task instances will run on, and the number of deleted rows is logged.
|
Hello @namanjain24-sudo - thank you for your contributions to Apache Airflow! The Airflow community has introduced a limit of 5 open pull requests at a time for contributors without write access to the repository. You currently have 7 open pull requests, so - as a one-time step of introducing the limit - we closed the ones where maintainers have not engaged yet:
These pull requests stay open because maintainers are already engaged in them - they count towards your limit:
This is not a judgement of you or of your changes. We never told contributors before that opening many pull requests at once was a problem, so there is nothing to feel bad about - and nothing is lost: your branches, commits and the review history stay where they are. What we ask you to do is to make your first prioritization decision: choose which of the pull requests above matter most to you, and reopen them (up to 5 open at a time, including the ones still open) with the "Reopen pull request" button or While your pull requests are waiting for review, the most valuable thing you can do is help in other ways - reviewing other contributors' pull requests, helping with issues, and taking part in the discussions on the devlist and Slack. Why we introduced the limit, what it means for you and how to reopen or restore a pull request is explained in https://github.com/apache/airflow/blob/main/contributing-docs/32_open_pull_request_limit.rst. Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting |
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 processeszwhile reading the checkpoint written fora. 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:
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 bothdag.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.rstdocuments 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.rstand in the clearing note inresumable-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:
test_cleartasks.pycover the cases in the table, a mapped task group, a second run of the same Dag that is left alone, anddag.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.{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.pypasses 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 andtest_dag.py. prek passes, including mypy for airflow-core.closes: #72725
Was generative AI tooling used to co-author this PR?
Generated-by: a Gen-AI coding assistant, following the guidelines. I reviewed the change and ran the checks above.