Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -203,7 +203,9 @@ Example skeleton:
class S3StateBackend(BaseStoreBackend):

def _task_ref(self, scope: TaskScope, key: str) -> str:
return f"airflow/task-store/{scope.dag_id}/{scope.run_id}/{scope.task_id}/{scope.map_index}/{key}"
# Two regions can hold the same task_id and map_index, so region_id must be part of the ref.
region = f"region={scope.region_id}/" if scope.region_id.int else ""
return f"airflow/task-store/{scope.dag_id}/{scope.run_id}/{region}{scope.task_id}/{scope.map_index}/{key}"

def _asset_ref(self, scope: AssetScope, key: str) -> str:
import hashlib
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,7 @@ Task and Asset state store provide two key/value stores to persist data like a j
- Default lifetime
- Primary use case
* - **Task state store**
- A single task Instance (``dag_id`` + ``run_id`` + ``task_id`` + ``map_index``)
- A single task Instance (``dag_id`` + ``run_id`` + ``task_id`` + ``region_id`` + ``map_index``)
- Configurable retention; cleared on task success when ``clear_on_success = True``
- Survive retries, track in-flight jobs, checkpoint progress within a run, resume progress from checkpoint set by a past run
* - **Asset state store**
Expand Down
4 changes: 3 additions & 1 deletion airflow-core/docs/core-concepts/task-state-store.rst
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,9 @@ Task State Store

.. versionadded:: 3.3

Task store is a persistent key/value store scoped to a single task instance (``dag_id`` + ``run_id`` + ``task_id`` + ``map_index``). It survives worker crashes and task retries within the same Dag run, making it suitable for storing external job IDs, intra-task checkpoints, and progress metadata.
Task store is a persistent key/value store scoped to a single task instance (``dag_id`` + ``run_id`` + ``task_id`` + ``region_id`` + ``map_index``). It survives worker crashes and task retries within the same Dag run, making it suitable for storing external job IDs, intra-task checkpoints, and progress metadata.
Comment thread
ashb marked this conversation as resolved.

``region_id`` is the all-zero UUID for task instances outside any dynamic region. Task instances created by dynamically expanded work, such as a loop, belong to a region, and two of them can share the same ``dag_id``, ``run_id``, ``task_id`` and ``map_index``. A custom backend must therefore include ``region_id`` in its uniqueness key.

Because it outlives a worker crash and stays readable by the next attempt, the task state store is the mechanism behind :ref:`durable execution <concepts-durable-execution>`, where a task continues
from where it stopped instead of repeating work or submitting a duplicate external job. Provider operators that advertise a ``durable`` parameter are built on the API described here.
Expand Down
4 changes: 3 additions & 1 deletion airflow-core/docs/migrations-ref.rst
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,9 @@ Here's the list of all the Database Migrations that are executed via when you ru
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| Revision ID | Revises ID | Airflow Version | Description |
+=========================+==================+===================+==============================================================+
| ``e7c2a91bd540`` (head) | ``90e4d18ccadf`` | ``3.4.0`` | Unify task attempt ownership without rewriting legacy XCom |
| ``54a27b6f9d01`` (head) | ``e7c2a91bd540`` | ``3.4.0`` | Add dynamic region storage. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``e7c2a91bd540`` | ``90e4d18ccadf`` | ``3.4.0`` | Unify task attempt ownership without rewriting legacy XCom |
| | | | data. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``90e4d18ccadf`` | ``e5a91c7f42b3`` | ``3.4.0`` | Add timetable_asset_gated to DagModel. |
Expand Down
72 changes: 72 additions & 0 deletions airflow-core/src/airflow/migrations/utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@

import contextlib
from contextlib import contextmanager
from textwrap import dedent


@contextmanager
Expand All @@ -30,6 +31,25 @@ def disable_sqlite_fkeys(op):
yield op


# Unlike disable_sqlite_fkeys, a failed SQLite upgrade rolls back atomically and foreign_keys is always restored.
@contextmanager
def sqlite_rebuilds(op):
if op.get_bind().dialect.name != "sqlite":
yield
return
if op.get_context().as_sql:
raise RuntimeError("SQLite offline SQL cannot render this migration's table rebuilds")
enabled = op.get_bind().exec_driver_sql("PRAGMA foreign_keys").scalar()
with op.get_context().autocommit_block():
op.execute("PRAGMA foreign_keys=OFF")
try:
with op.get_bind().begin_nested():
yield
finally:
with op.get_context().autocommit_block():
op.execute(f"PRAGMA foreign_keys={int(enabled)}")


def mysql_drop_foreignkey_if_exists(constraint_name, table_name, op):
"""Older Mysql versions do not support DROP FOREIGN KEY IF EXISTS."""
op.execute(f"""
Expand All @@ -55,6 +75,58 @@ def mysql_drop_foreignkey_if_exists(constraint_name, table_name, op):
""")


def raise_if_rows_exist(select_sql: str, message: str, op) -> None:
"""
Abort the migration when ``select_sql`` returns a row, using SQL so it also runs offline.

Call it before any DDL: MySQL DDL is not transactional, so a late failure leaves a half-applied
migration. ``message`` must not contain single quotes or percent signs and, for MySQL, must be at
most 128 characters. PostgreSQL and MySQL only.
"""
dialect = op.get_context().dialect.name
if dialect == "postgresql":
op.execute(
dedent("""
DO $$
BEGIN
IF EXISTS ({select_sql}) THEN
RAISE EXCEPTION '{message}';
END IF;
END $$
""").format(select_sql=select_sql, message=message)
)
elif dialect == "mysql":
create = (
dedent("""
CREATE PROCEDURE AirflowMigrationGuard()
BEGIN
IF EXISTS ({select_sql}) THEN
SIGNAL SQLSTATE '45000' SET MESSAGE_TEXT = '{message}';
END IF;
END
""")
.strip()
.format(select_sql=select_sql, message=message)
)
context = op.get_context()
if context.as_sql:
# The mysql client splits on ';' unless the procedure body is fenced by DELIMITER.
context.output_buffer.write(
"DROP PROCEDURE IF EXISTS AirflowMigrationGuard;\n"
f"DELIMITER //\n{create}//\nDELIMITER ;\n"
"CALL AirflowMigrationGuard();\nDROP PROCEDURE AirflowMigrationGuard;\n\n"
)
else:
op.execute("DROP PROCEDURE IF EXISTS AirflowMigrationGuard")
op.execute(create)
try:
op.execute("CALL AirflowMigrationGuard()")
finally:
op.execute("DROP PROCEDURE AirflowMigrationGuard")
else:
raise NotImplementedError(f"No SQL guard for dialect {dialect}")


def ignore_sqlite_value_error():
from alembic import op

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,13 +26,12 @@

from __future__ import annotations

from contextlib import contextmanager

import sqlalchemy as sa
from alembic import op
from sqlalchemy.dialects import postgresql, sqlite

from airflow.migrations.db_types import TIMESTAMP, StringID
from airflow.migrations.utils import sqlite_rebuilds
from airflow.utils.sqlalchemy import ExecutorConfigType, ExtendedJSON, UtcDateTime

revision = "e7c2a91bd540"
Expand Down Expand Up @@ -99,25 +98,6 @@
)


# Unlike disable_sqlite_fkeys, a failed SQLite upgrade rolls back atomically and foreign_keys is always restored.
@contextmanager
def _sqlite_rebuilds():
if op.get_bind().dialect.name != "sqlite":
yield
return
if op.get_context().as_sql:
raise RuntimeError("SQLite offline SQL cannot render this migration's table rebuilds")
enabled = op.get_bind().exec_driver_sql("PRAGMA foreign_keys").scalar()
with op.get_context().autocommit_block():
op.execute("PRAGMA foreign_keys=OFF")
try:
with op.get_bind().begin_nested():
yield
finally:
with op.get_context().autocommit_block():
op.execute(f"PRAGMA foreign_keys={int(enabled)}")


def _check_source():
bind = op.get_bind()
if op.get_context().as_sql:
Expand Down Expand Up @@ -190,7 +170,7 @@ def _redirect_legacy(table, constraint, target, *, onupdate=None, not_valid=Fals
def upgrade():
"""Retain attempts and give legacy and new task data immutable UUID owners."""
_check_source()
with _sqlite_rebuilds():
with sqlite_rebuilds(op):
owner = op.create_table(
"legacy_task_data_owner",
sa.Column("dag_id", StringID(), nullable=False),
Expand Down Expand Up @@ -387,7 +367,7 @@ def downgrade():
)
):
raise RuntimeError("Cannot downgrade attempt ownership after legacy owners changed coordinates")
with _sqlite_rebuilds():
with sqlite_rebuilds(op):
op.drop_table("rtif_v2")
op.drop_table("xcom_v2")
op.drop_index("ti_current_state", table_name="task_instance")
Expand Down
Loading
Loading