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],

Reply via email to