Skip to content

Introduce the (internal/db level) concept of "dynamic regions" for TIs - #74275

Merged
ashb merged 1 commit into
task-loops-stack-1from
task-loops-stack-2
Oct 10, 2026
Merged

ashb merged 1 commit into
task-loops-stack-1from
task-loops-stack-2

Conversation

@ashb

@ashb ashb commented Oct 5, 2026 •

Copy link
Copy Markdown
Member

As mentioned in the first docs PR of this stack, and the precursor PR #74222,
we need a way of distinguishing between "interation 2" and "iteration 2 after
we already cleared and re-ran loop 2 but-these-are-totally-different-TIs, not
retries.", or more simply in the digram as 2' (said as "two prime").

Task loops need the same task_id to exist once per pass, and a clear needs to
put replacement task beside the work it has "retired" or superceded. A task
instance is identified by dag_id, task_id, run_id and map_index today, so
neither is possible: the second row would collide with the first as they both
have the same task_id and try_number.

A "dynamic region is one execution of a construct that creates task instances at
run time: a loop, or the expansion of a mapped task or task group. Each region
is a row recording which Dag run it belongs to, which node it executes (node_id,
the id of the task group or mapped task), which enclosing region and position
contains it and (so that we can have mapped tasks inside loops!), when a clear
replaced an earlier execution, which region it was forked from and where it
resumes (i.e. "clear iteration 2 onwards").

A region is a row rather than a position (i.e. it's own table, not just a colum
on TI) because a position can be reused/"overwritten" after a clear, and would
then stop saying which precise TI created the data. Regions are immutable for
the same reason; everything that changes stays on the task instances.

This change only adds the storage to hopefully make PR creation easier.
Nothing creates a region yet, and existing rows take the all-zero sentinel,
meaning outside any region, because rewriting existing rows on upgrade was
measured at production scale and rejected (4+ hrs! But I never let it run to
completion. Even 30mins is too long). XCom and rendered fields are keyed by
attempt UUID since the precursor ownership PR and so don't need a region column.
region_index is an ORM synonym of map_index for now.

Later changes will build upon this to make mapped expansions and loops create
regions, resolve and write task data by (region_id, region_index), add the task
loop gate that decides whether a loop runs another pass/iteration, and let a
clear replace only the selected work by forking a region instead of rewriting
history.

Downgrade refuses when rows differ only by region, because dropping the column
would otherwise merge distinct executions. The check scans task_instance, so
expect it to be slow on very large installations.


Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

  • Read the Pull Request Guidelines for more information. Note: commit author/co-author name and email in commits become permanently public when merged.
  • For fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
  • When adding dependency, check compliance with the ASF 3rd Party License Policy.
  • For significant user-facing changes create newsfragment: {pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.

@ashb
ashb added this pull request to stack #74276 October 5, 2026 16:35
@ashb
ashb removed this pull request from stack #74276 October 5, 2026 16:36
@ashb
ashb added this pull request to stack #74339 October 6, 2026 13:52
@ashb
ashb force-pushed the task-loops-stack-2 branch from 67a2c41 to 58c9415 Compare October 6, 2026 14:55
@ashb
ashb marked this pull request as ready for review October 6, 2026 15:40
Comment thread shared/state/src/airflow_shared/state/__init__.py
Comment thread airflow-core/tests/unit/migrations/test_0143_add_dynamic_region_storage.py Outdated
Comment thread airflow-core/tests/unit/migrations/test_0143_add_dynamic_region_storage.py Outdated
Comment thread airflow-core/tests/unit/models/test_taskinstance.py Outdated
Comment thread airflow-core/src/airflow/models/taskinstance.py Outdated
Comment thread airflow-core/src/airflow/models/task_state_store.py
@ashb
ashb force-pushed the task-loops-stack-2 branch from 58c9415 to 5678908 Compare October 7, 2026 13:57
@ashb
ashb removed this pull request from stack #74339 October 7, 2026 15:20
@ashb
ashb force-pushed the task-loops-stack-2 branch from 5678908 to 2299adf Compare October 7, 2026 15:22
@ashb
ashb added this pull request to stack #74410 October 7, 2026 15:22
@ashb
ashb force-pushed the task-loops-stack-2 branch 2 times, most recently from 9fbf5c3 to 8bca741 Compare October 8, 2026 13:48
@ashb
ashb force-pushed the task-loops-stack-2 branch from 8bca741 to 049690d Compare October 8, 2026 13:53
@ashb
ashb force-pushed the task-loops-stack-2 branch from 049690d to 8dd908b Compare October 8, 2026 16:06
Comment thread airflow-core/src/airflow/utils/sqlalchemy.py Outdated
Comment thread airflow-core/docs/core-concepts/task-state-store.rst
@ashb
ashb force-pushed the task-loops-stack-2 branch from 8dd908b to fe9cd34 Compare October 8, 2026 20:45
@ashb
ashb force-pushed the task-loops-stack-2 branch 2 times, most recently from 6cab9a2 to a6f0f2c Compare October 9, 2026 14:03
@ashb
ashb force-pushed the task-loops-stack-2 branch from a6f0f2c to 205af7f Compare October 9, 2026 16:28
@ashb
ashb force-pushed the task-loops-stack-2 branch from 205af7f to 9a110d1 Compare October 9, 2026 20:34
@ashb
ashb force-pushed the task-loops-stack-2 branch 2 times, most recently from 30c3b24 to e1d3269 Compare October 9, 2026 22:18
As mentioned in the first docs PR of this stack, and the precursor PR #74222,
we need a way of distinguishing between "interation 2" and "iteration 2 after
we already cleared and re-ran loop 2 but-these-are-totally-different-TIs, not
retries.", or more simply in the digram as `2'` (said as "two prime").

Task loops need the same task_id to exist once per pass, and a clear needs to
put replacement task beside the work it has "retired" or superceded. A task
instance is identified by dag_id, task_id, run_id and map_index today, so
neither is possible: the second row would collide with the first as they both
have the same task_id and try_number.

A "dynamic region is one execution of a construct that creates task instances at
run time: a loop, or the expansion of a mapped task or task group. Each region
is a row recording which Dag run it belongs to, which node it executes (node_id,
the id of the task group or mapped task), which enclosing region and position
contains it and (so that we can have mapped tasks inside loops!), when a clear
replaced an earlier execution, which region it was forked from and where it
resumes (i.e. "clear iteration 2 onwards").

A region is a row rather than a position (i.e. it's own table, not just a colum
on TI) because a position can be reused/"overwritten" after a clear, and would
then stop saying which precise TI created the data. Regions are immutable for
the same reason; everything that changes stays on the task instances.

This change only adds the storage to hopefully make PR creation easier.
Nothing creates a region yet, and existing rows take the all-zero sentinel,
meaning outside any region, because rewriting existing rows on upgrade was
measured at production scale and rejected (4+ hrs! But I never let it run to
completion. Even 30mins is too long). XCom and rendered fields are keyed by
attempt UUID since the precursor ownership PR and so don't need a region column.
`region_index` is an ORM synonym of `map_index` for now.

Later changes will build upon this to make mapped expansions and loops create
regions, resolve and write task data by (region_id, region_index), add the task
loop gate that decides whether a loop runs another pass/iteration, and let a
clear replace only the selected work by forking a region instead of rewriting
history.

Downgrade refuses when rows differ only by region, because dropping the column
would otherwise merge distinct executions. The check scans task_instance, so
expect it to be slow on very large installations.
@ashb
ashb force-pushed the task-loops-stack-2 branch from e1d3269 to cdda9db Compare October 10, 2026 07:05
@ashb
ashb merged commit 78572f6 into main Oct 10, 2026
161 of 167 checks passed
@ashb
ashb deleted the task-loops-stack-2 branch October 10, 2026 22:49
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants