Skip to content

Attribute log, asset-event and asset-state writes to the task attempt - #74347

Open
ashb wants to merge 1 commit into
task-loops-stack-5from
task-loops-stack-6
Open

ashb wants to merge 1 commit into
task-loops-stack-5from
task-loops-stack-6

Conversation

@ashb

@ashb ashb commented Oct 6, 2026 •

Copy link
Copy Markdown
Member

Audit log rows, asset events and asset-state writers identified their task by
(dag_id, task_id, run_id, map_index). Once a loop pass or a clear can reuse
those coordinates for another execution, a row read back by coordinates points
at the wrong attempt. Retired attempts stay in task_instance under their own
UUID, so the UUID is the strongest anchor: a row resolves to the attempt that
actually acted, and keeps doing so after that attempt is retired.

Rows written before the upgrade only have coordinates. We leave them without
an attributed attempt instead of matching them to a guess, because a wrong
attribution is worse than none.

An attributed row's public map index comes from the try's pinned definition,
as it does everywhere else, and is NULL once the attempt row has been purged
because nothing is left to derive it from. The asset-state writer also records
its region id, region index and try number so it stays describable after cleanup
removes the attempt. The stuck-in-queued accounting counts log rows by attempt
UUID and falls back to coordinates only for sentinel-region rows, which are the
rows that can still be unambiguous by coordinates.

OpenLineage uses the attempt UUID as the run id of regional tasks so the
listener events, the lineage macros and downstream asset dependencies agree on
which run they describe.

Task notes need no new storage: a note hangs off the attempt UUID and a retired
attempt keeps its row, so the retain-notes migration from the earlier design is
not carried over.


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 #74339 October 6, 2026 15:06
@boring-cyborg boring-cyborg Bot added area:airflow-ctl area:API Airflow's REST/HTTP API area:db-migrations PRs with DB migration area:Executors-core LocalExecutor & SequentialExecutor area:providers area:scheduler area:task-sdk area:UI Related to UI/UX. For Frontend Developers. kind:documentation provider:openlineage AIP-53 labels Oct 6, 2026
Comment thread airflow-core/src/airflow/api_fastapi/core_api/services/public/event_logs.py Outdated
Comment thread airflow-core/src/airflow/models/log.py
Comment thread providers/openlineage/src/airflow/providers/openlineage/utils/utils.py Outdated
Comment thread providers/openlineage/src/airflow/providers/openlineage/utils/utils.py Outdated
Comment thread airflow-core/src/airflow/api_fastapi/core_api/routes/public/event_logs.py Outdated
@ashb
ashb force-pushed the task-loops-stack-6 branch from d25a7c8 to 2e69b13 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-6 branch from 2e69b13 to 7be9b12 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-6 branch 2 times, most recently from 8a05c82 to 6bb15d6 Compare October 7, 2026 20:57
@ashb
ashb force-pushed the task-loops-stack-6 branch from 6bb15d6 to 5c17d4d Compare October 8, 2026 13:48
@ashb
ashb force-pushed the task-loops-stack-6 branch 2 times, most recently from 0e415ae to c3b6ed7 Compare October 8, 2026 16:06
Comment thread airflow-core/src/airflow/models/log.py
Comment thread providers/openlineage/src/airflow/providers/openlineage/utils/utils.py Outdated
Comment thread airflow-core/src/airflow/api_fastapi/core_api/services/public/event_logs.py Outdated
Comment thread providers/openlineage/src/airflow/providers/openlineage/utils/utils.py Outdated
Comment thread providers/openlineage/tests/unit/openlineage/utils/test_utils.py Outdated
@ashb
ashb force-pushed the task-loops-stack-6 branch from c3b6ed7 to 8d6022b Compare October 8, 2026 20:45
Audit log rows, asset events and asset-state writers identified their task by
(dag_id, task_id, run_id, map_index). Once a loop pass or a clear can reuse
those coordinates for another execution, a row read back by coordinates points
at the wrong attempt. Retired attempts stay in task_instance under their own
UUID, so the UUID is the strongest anchor: a row resolves to the attempt that
actually acted, and keeps doing so after that attempt is retired.

Rows written before the upgrade only have coordinates. We leave them without
an attributed attempt instead of matching them to a guess, because a wrong
attribution is worse than none.

An attributed row's public map index comes from the attempt's pinned definition,
as it does everywhere else, and is NULL once the attempt row has been purged
because nothing is left to derive it from. The asset-state writer also records
its region id, region index and try number so it stays describable after cleanup
removes the attempt. The stuck-in-queued accounting counts log rows by attempt
UUID and falls back to coordinates only for sentinel-region rows, which are the
rows that can still be unambiguous by coordinates.

OpenLineage uses the attempt UUID as the run id of regional tasks so the
listener events, the lineage macros and downstream asset dependencies agree on
which run they describe.

Task notes need no new storage: a note hangs off the attempt UUID and a retired
attempt keeps its row, so the retain-notes migration from the earlier design is
not carried over.
@ashb
ashb force-pushed the task-loops-stack-6 branch from 0f448fb to 094f4b6 Compare October 10, 2026 07:05

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:airflow-ctl area:API Airflow's REST/HTTP API area:db-migrations PRs with DB migration area:Executors-core LocalExecutor & SequentialExecutor area:providers area:scheduler area:task-sdk area:UI Related to UI/UX. For Frontend Developers. kind:documentation provider:openlineage AIP-53

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants