This is an automated email from the ASF dual-hosted git repository. ashb pushed a commit to branch xcom-mapped-length-check-not-valid in repository https://gitbox.apache.org/repos/asf/airflow.git
commit 56a8bee24f2c100776d07fd88652e7c53cb95c16 Author: Ash Berlin-Taylor <[email protected]> AuthorDate: Wed Oct 7 16:09:15 2026 +0100 Speed up Migration on 0136 It turns out that adding a CHECK constraint on mysql _copies the entire table_. Large values and all. Enforcing this check is simply not worth a possibly 4 hour migration! --- .../0136_3_4_0_fold_task_map_into_xcom_mapped_length.py | 12 ++++++++++-- airflow-core/src/airflow/models/xcom.py | 11 +++++++++-- 2 files changed, 19 insertions(+), 4 deletions(-) diff --git a/airflow-core/src/airflow/migrations/versions/0136_3_4_0_fold_task_map_into_xcom_mapped_length.py b/airflow-core/src/airflow/migrations/versions/0136_3_4_0_fold_task_map_into_xcom_mapped_length.py index 86e1e7ccf4a..87d3b2421e2 100644 --- a/airflow-core/src/airflow/migrations/versions/0136_3_4_0_fold_task_map_into_xcom_mapped_length.py +++ b/airflow-core/src/airflow/migrations/versions/0136_3_4_0_fold_task_map_into_xcom_mapped_length.py @@ -116,9 +116,16 @@ def build_restore_statement(xcom_name: str = "xcom", task_map_name: str = "task_ def upgrade(): """Fold task_map into xcom.mapped_length.""" with disable_sqlite_fkeys(op): + dialect = op.get_bind().dialect.name with op.batch_alter_table("xcom", schema=None) as batch_op: batch_op.add_column(sa.Column("mapped_length", sa.Integer(), nullable=True)) - batch_op.create_check_constraint("mapped_length_not_negative", "mapped_length >= 0") + # MySQL can only add a CHECK by copying the whole table (ALGORITHM=COPY). + if dialect != "mysql": + batch_op.create_check_constraint( + "mapped_length_not_negative", + "mapped_length >= 0", + postgresql_not_valid=True, + ) op.execute(build_backfill_statement()) op.drop_table("task_map") @@ -154,5 +161,6 @@ def downgrade(): op.execute(build_restore_statement()) with op.batch_alter_table("xcom", schema=None) as batch_op: - batch_op.drop_constraint("mapped_length_not_negative", type_="check") + if op.get_bind().dialect.name != "mysql": + batch_op.drop_constraint("mapped_length_not_negative", type_="check") batch_op.drop_column("mapped_length") diff --git a/airflow-core/src/airflow/models/xcom.py b/airflow-core/src/airflow/models/xcom.py index e64789cb8ce..a2cc63c4413 100644 --- a/airflow-core/src/airflow/models/xcom.py +++ b/airflow-core/src/airflow/models/xcom.py @@ -60,7 +60,7 @@ from airflow.utils.sqlalchemy import UtcDateTime, build_upsert_stmt log = logging.getLogger(__name__) if TYPE_CHECKING: - from sqlalchemy.engine import Row + from sqlalchemy.engine import Dialect, Row from sqlalchemy.orm import Session from sqlalchemy.orm.util import AliasedClass from sqlalchemy.sql.expression import CompoundSelect, Select, Subquery, TextClause @@ -72,6 +72,11 @@ if TYPE_CHECKING: XCOM_RETURN_KEY = "return_value" +def _is_not_mysql(*args: Any, dialect: Dialect, **kwargs: Any) -> bool: + # MySQL can only add this CHECK by copying the whole table, so migration 0136 does not create it there. + return dialect.name != "mysql" + + class XComModelV1(TaskInstanceDependencies): """XCom values stored before Airflow 3.4, keyed by dag, task, run and map index.""" @@ -99,7 +104,9 @@ class XComModelV1(TaskInstanceDependencies): # separately, and enforce uniqueness with DagRun.id instead. Index("idx_xcom_key", key), Index("idx_xcom_task_instance", dag_id, task_id, run_id, map_index), - CheckConstraint(mapped_length >= 0, name=conv("ck_xcom_mapped_length_not_negative")), + CheckConstraint(mapped_length >= 0, name=conv("ck_xcom_mapped_length_not_negative")).ddl_if( + callable_=_is_not_mysql + ), PrimaryKeyConstraint("dag_run_id", "task_id", "map_index", "key", name="xcom_pkey"), ForeignKeyConstraint( [dag_id, task_id, run_id, map_index],
