From 2790cb25b328196e79a0b1ff31575296bb34730b Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Thu, 6 Aug 2026 12:05:41 -0700 Subject: [PATCH 1/5] Allow configuring the collation of asset name columns The asset name, uri and group columns hard-code the latin1_general_cs collation on MySQL. Several MySQL-compatible engines do not provide that collation, so Airflow cannot create its own schema on them even though the rest of the database works. There is no way to override it from outside, because the collation is baked into the ORM column definitions. closes: #31373 Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../src/airflow/config_templates/config.yml | 11 ++ airflow-core/src/airflow/models/asset.py | 103 ++---------------- airflow-core/src/airflow/models/base.py | 17 +++ 3 files changed, 39 insertions(+), 92 deletions(-) diff --git a/airflow-core/src/airflow/config_templates/config.yml b/airflow-core/src/airflow/config_templates/config.yml index 49ae0c1f3090d..21a27bcfb815a 100644 --- a/airflow-core/src/airflow/config_templates/config.yml +++ b/airflow-core/src/airflow/config_templates/config.yml @@ -649,6 +649,17 @@ database: type: string example: ~ default: ~ + sql_engine_collation_for_asset_names: + description: | + Collation for the ``name``, ``uri`` and ``group`` columns of the asset tables on + ``mysql`` and ``mariadb``. These columns hold ASCII values and are indexed at 1500 + characters, so a single-byte charset is used to stay within the maximum index size. + Override this if your database engine does not provide ``latin1_general_cs`` -- + for example TiDB, which supports ``latin1_bin`` instead. + version_added: 3.4.0 + type: string + example: "latin1_bin" + default: "latin1_general_cs" sql_alchemy_pool_enabled: description: | If SQLAlchemy should pool database connections. diff --git a/airflow-core/src/airflow/models/asset.py b/airflow-core/src/airflow/models/asset.py index d980c2d69acec..2076946de3d62 100644 --- a/airflow-core/src/airflow/models/asset.py +++ b/airflow-core/src/airflow/models/asset.py @@ -30,7 +30,6 @@ Index, Integer, PrimaryKeyConstraint, - String, Table, delete, select, @@ -40,7 +39,7 @@ from airflow._shared.timezones import timezone from airflow.configuration import conf as airflow_conf -from airflow.models.base import Base, StringID +from airflow.models.base import ASSET_STR_FIELD, Base, StringID from airflow.utils.sqlalchemy import UtcDateTime if TYPE_CHECKING: @@ -148,15 +147,7 @@ class AssetWatcherModel(Base): """A table to store asset watchers.""" name: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, nullable=False, ) asset_id: Mapped[int] = mapped_column(Integer, primary_key=True, nullable=False) @@ -199,27 +190,11 @@ class AssetAliasModel(Base): id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) name: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, nullable=False, ) group: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, default="", nullable=False, ) @@ -280,39 +255,15 @@ class AssetModel(Base): id: Mapped[int] = mapped_column(Integer, primary_key=True, autoincrement=True) name: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, nullable=False, ) uri: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, nullable=False, ) group: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, default=str, nullable=False, ) @@ -409,27 +360,11 @@ class AssetActive(Base): """ name: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, nullable=False, ) uri: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, nullable=False, ) @@ -457,15 +392,7 @@ class DagScheduleAssetNameReference(Base): """Reference from a DAG to an asset name reference of which it is a consumer.""" name: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, primary_key=True, nullable=False, ) @@ -503,15 +430,7 @@ class DagScheduleAssetUriReference(Base): """Reference from a DAG to an asset URI reference of which it is a consumer.""" uri: Mapped[str] = mapped_column( - String(length=1500).with_variant( - String( - length=1500, - # latin1 allows for more indexed length in mysql - # and this field should only be ascii chars - collation="latin1_general_cs", - ), - "mysql", - ), + ASSET_STR_FIELD, primary_key=True, nullable=False, ) diff --git a/airflow-core/src/airflow/models/base.py b/airflow-core/src/airflow/models/base.py index 1c7af7b275ab3..2840bf7e8f8a8 100644 --- a/airflow-core/src/airflow/models/base.py +++ b/airflow-core/src/airflow/models/base.py @@ -83,6 +83,23 @@ def get_id_collation_args(): COLLATION_ARGS: dict[str, Any] = get_id_collation_args() +def get_asset_str_field(length: int = 1500) -> String: + """ + Build the string type used for asset name/uri/group columns. + + On MySQL these carry an explicit latin1 collation: the values are ASCII, and + a 1-byte-per-character charset keeps the 1500-char unique indexes inside the + 3072-byte index limit that utf8mb4 would blow past. The collation is + overridable because MySQL-compatible engines do not all ship + ``latin1_general_cs`` (TiDB, for one, accepts only ``latin1_bin``). + """ + collation = conf.get("database", "sql_engine_collation_for_asset_names", fallback="latin1_general_cs") + return String(length=length).with_variant(String(length=length, collation=collation), "mysql") + + +ASSET_STR_FIELD: String = get_asset_str_field() + + def StringID(*, length=ID_LEN, **kwargs) -> String: return String(length=length, **kwargs, **COLLATION_ARGS) From 86dfbd019148aa3473cf55a6831aa4b87ddaf857 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Fri, 7 Aug 2026 12:20:22 -0700 Subject: [PATCH 2/5] Apply the asset collation setting to migrations as well as models The asset name/uri/group columns carry an explicit MySQL collation so their 1500-character unique indexes stay inside the 3072-byte index limit. That collation is configurable on the models, but the migrations that create and alter those columns still hard-coded `latin1_general_cs`, so the setting only took effect when the schema was created from the ORM. A fresh install on a MySQL-compatible engine that lacks that collation replays the migrations instead and fails on the first asset table, which leaves the setting useless in exactly the case it was added for. Resolve the collation at migration run time through the same configuration key the models read. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- airflow-core/src/airflow/migrations/utils.py | 13 +++++++++++++ .../versions/0000_2_6_2_squashed_migrations.py | 3 ++- .../versions/0022_2_10_0_add_dataset_alias.py | 4 +++- .../0036_3_0_0_add_name_field_to_dataset_model.py | 6 ++++-- .../versions/0038_3_0_0_add_asset_active.py | 4 +++- ...39_3_0_0_tweak_assetaliasmodel_to_match_asset.py | 6 ++++-- .../0054_3_0_0_add_asset_reference_models.py | 3 ++- ..._3_2_0_replace_asset_trigger_table_with_asset.py | 4 +++- 8 files changed, 34 insertions(+), 9 deletions(-) diff --git a/airflow-core/src/airflow/migrations/utils.py b/airflow-core/src/airflow/migrations/utils.py index 4eeaf373c6a87..226928ca3b843 100644 --- a/airflow-core/src/airflow/migrations/utils.py +++ b/airflow-core/src/airflow/migrations/utils.py @@ -61,3 +61,16 @@ def ignore_sqlite_value_error(): if op.get_bind().dialect.name == "sqlite": return contextlib.suppress(ValueError) return contextlib.nullcontext() + + +def asset_name_collation() -> str: + """ + Return the MySQL collation for asset name/uri/group columns. + + Mirrors ``airflow.models.base.get_asset_str_field``. Migrations resolve it at + run time rather than hard-coding it, because MySQL-compatible engines do not + all ship ``latin1_general_cs`` and a fresh install replays these migrations. + """ + from airflow.configuration import conf + + return conf.get("database", "sql_engine_collation_for_asset_names", fallback="latin1_general_cs") diff --git a/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py b/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py index 739e4e7f88768..49f30872754d0 100644 --- a/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py +++ b/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py @@ -34,6 +34,7 @@ from sqlalchemy.dialects.mysql import MEDIUMTEXT from airflow.migrations.db_types import StringID +from airflow.migrations.utils import asset_name_collation from airflow.utils.sqlalchemy import ExtendedJSON, UtcDateTime # revision identifiers, used by Alembic. @@ -208,7 +209,7 @@ def upgrade() -> None: length=3000, # latin1 allows for more indexed length in mysql # and this field should only be ascii chars - collation="latin1_general_cs", + collation=asset_name_collation(), ), "mysql", ), diff --git a/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py b/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py index 0d4a9efe0ebac..0cbe6841a4f62 100644 --- a/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py +++ b/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py @@ -30,6 +30,8 @@ import sqlalchemy as sa from alembic import op +from airflow.migrations.utils import asset_name_collation + # revision identifiers, used by Alembic. revision = "05e19f3176be" down_revision = "d482b7261ff9" @@ -46,7 +48,7 @@ def upgrade(): sa.Column( "name", sa.String(length=3000).with_variant( - sa.String(length=3000, collation="latin1_general_cs"), "mysql" + sa.String(length=3000, collation=asset_name_collation()), "mysql" ), nullable=False, ), diff --git a/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py b/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py index b1f925dbffca2..a5a5cb5512b95 100644 --- a/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py +++ b/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py @@ -39,6 +39,8 @@ import sqlalchemy as sa from alembic import op +from airflow.migrations.utils import asset_name_collation + # revision identifiers, used by Alembic. revision = "0d9e73a75ee4" down_revision = "44eabb1904b4" @@ -47,7 +49,7 @@ airflow_version = "3.0.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation="latin1_general_cs"), + sa.String(length=1500, collation=asset_name_collation()), "mysql", ) @@ -127,7 +129,7 @@ def downgrade(): batch_op.alter_column( "uri", type_=sa.String(length=3000).with_variant( - sa.String(length=3000, collation="latin1_general_cs"), + sa.String(length=3000, collation=asset_name_collation()), "mysql", ), nullable=False, diff --git a/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py b/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py index 2a992cab4126e..39b8d8151692b 100644 --- a/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py +++ b/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py @@ -30,6 +30,8 @@ import sqlalchemy as sa from alembic import op +from airflow.migrations.utils import asset_name_collation + # revision identifiers, used by Alembic. revision = "5a5d66100783" down_revision = "c3389cd7793f" @@ -38,7 +40,7 @@ airflow_version = "3.0.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation="latin1_general_cs"), + sa.String(length=1500, collation=asset_name_collation()), "mysql", ) diff --git a/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py b/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py index d0067f1288255..79ded55d6d83d 100644 --- a/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py +++ b/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py @@ -42,6 +42,8 @@ import sqlalchemy as sa from alembic import op +from airflow.migrations.utils import asset_name_collation + # Revision identifiers, used by Alembic. revision = "fb2d4922cd79" down_revision = "5a5d66100783" @@ -50,7 +52,7 @@ airflow_version = "3.0.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation="latin1_general_cs"), + sa.String(length=1500, collation=asset_name_collation()), "mysql", ) @@ -76,7 +78,7 @@ def downgrade(): batch_op.alter_column( "name", type_=sa.String(length=3000).with_variant( - sa.String(length=3000, collation="latin1_general_cs"), + sa.String(length=3000, collation=asset_name_collation()), "mysql", ), nullable=False, diff --git a/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py b/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py index 4a34da03ecc2a..0eb50af20adb7 100644 --- a/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py +++ b/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py @@ -30,6 +30,7 @@ from alembic import op from airflow.migrations.db_types import StringID +from airflow.migrations.utils import asset_name_collation from airflow.utils.sqlalchemy import UtcDateTime # revision identifiers, used by Alembic. @@ -40,7 +41,7 @@ airflow_version = "3.0.0" ASSET_STR_FIELD = sa.String(length=1500).with_variant( - sa.String(length=1500, collation="latin1_general_cs"), "mysql" + sa.String(length=1500, collation=asset_name_collation()), "mysql" ) diff --git a/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py b/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py index 81c134741f4d2..17ab84969d16a 100644 --- a/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py +++ b/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py @@ -30,6 +30,8 @@ import sqlalchemy as sa from alembic import op +from airflow.migrations.utils import asset_name_collation + # revision identifiers, used by Alembic. revision = "15d84ca19038" down_revision = "509b94a1042d" @@ -38,7 +40,7 @@ airflow_version = "3.2.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation="latin1_general_cs"), + sa.String(length=1500, collation=asset_name_collation()), "mysql", ) From 12853319e4e409b745494415f365d1f2805d60b3 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Mon, 24 Aug 2026 01:24:57 -0400 Subject: [PATCH 3/5] Validate configurable asset collations The setting controls both ORM table creation and historical migrations, so both paths need focused coverage to keep unsupported hard-coded collations from returning. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- airflow-core/src/airflow/migrations/utils.py | 14 ++++---------- .../versions/0000_2_6_2_squashed_migrations.py | 4 ++-- .../versions/0022_2_10_0_add_dataset_alias.py | 4 ++-- ...0036_3_0_0_add_name_field_to_dataset_model.py | 6 +++--- .../versions/0038_3_0_0_add_asset_active.py | 4 ++-- ...3_0_0_tweak_assetaliasmodel_to_match_asset.py | 6 +++--- .../0054_3_0_0_add_asset_reference_models.py | 4 ++-- ...2_0_replace_asset_trigger_table_with_asset.py | 4 ++-- airflow-core/src/airflow/models/base.py | 10 +--------- .../unit/migrations/test_migration_utils.py | 9 +++++++++ airflow-core/tests/unit/models/test_base.py | 16 +++++++++++++++- 11 files changed, 45 insertions(+), 36 deletions(-) diff --git a/airflow-core/src/airflow/migrations/utils.py b/airflow-core/src/airflow/migrations/utils.py index 226928ca3b843..4fcd22cf3c262 100644 --- a/airflow-core/src/airflow/migrations/utils.py +++ b/airflow-core/src/airflow/migrations/utils.py @@ -19,6 +19,8 @@ import contextlib from contextlib import contextmanager +from airflow.configuration import conf + @contextmanager def disable_sqlite_fkeys(op): @@ -63,14 +65,6 @@ def ignore_sqlite_value_error(): return contextlib.nullcontext() -def asset_name_collation() -> str: - """ - Return the MySQL collation for asset name/uri/group columns. - - Mirrors ``airflow.models.base.get_asset_str_field``. Migrations resolve it at - run time rather than hard-coding it, because MySQL-compatible engines do not - all ship ``latin1_general_cs`` and a fresh install replays these migrations. - """ - from airflow.configuration import conf - +def get_asset_name_collation() -> str: + """Return the MySQL collation used by asset identifier columns.""" return conf.get("database", "sql_engine_collation_for_asset_names", fallback="latin1_general_cs") diff --git a/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py b/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py index 49f30872754d0..a1cca6f7f6f4e 100644 --- a/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py +++ b/airflow-core/src/airflow/migrations/versions/0000_2_6_2_squashed_migrations.py @@ -34,7 +34,7 @@ from sqlalchemy.dialects.mysql import MEDIUMTEXT from airflow.migrations.db_types import StringID -from airflow.migrations.utils import asset_name_collation +from airflow.migrations.utils import get_asset_name_collation from airflow.utils.sqlalchemy import ExtendedJSON, UtcDateTime # revision identifiers, used by Alembic. @@ -209,7 +209,7 @@ def upgrade() -> None: length=3000, # latin1 allows for more indexed length in mysql # and this field should only be ascii chars - collation=asset_name_collation(), + collation=get_asset_name_collation(), ), "mysql", ), diff --git a/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py b/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py index 0cbe6841a4f62..7d4f4a970c0f2 100644 --- a/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py +++ b/airflow-core/src/airflow/migrations/versions/0022_2_10_0_add_dataset_alias.py @@ -30,7 +30,7 @@ import sqlalchemy as sa from alembic import op -from airflow.migrations.utils import asset_name_collation +from airflow.migrations.utils import get_asset_name_collation # revision identifiers, used by Alembic. revision = "05e19f3176be" @@ -48,7 +48,7 @@ def upgrade(): sa.Column( "name", sa.String(length=3000).with_variant( - sa.String(length=3000, collation=asset_name_collation()), "mysql" + sa.String(length=3000, collation=get_asset_name_collation()), "mysql" ), nullable=False, ), diff --git a/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py b/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py index a5a5cb5512b95..53e4efb9d2704 100644 --- a/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py +++ b/airflow-core/src/airflow/migrations/versions/0036_3_0_0_add_name_field_to_dataset_model.py @@ -39,7 +39,7 @@ import sqlalchemy as sa from alembic import op -from airflow.migrations.utils import asset_name_collation +from airflow.migrations.utils import get_asset_name_collation # revision identifiers, used by Alembic. revision = "0d9e73a75ee4" @@ -49,7 +49,7 @@ airflow_version = "3.0.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation=asset_name_collation()), + sa.String(length=1500, collation=get_asset_name_collation()), "mysql", ) @@ -129,7 +129,7 @@ def downgrade(): batch_op.alter_column( "uri", type_=sa.String(length=3000).with_variant( - sa.String(length=3000, collation=asset_name_collation()), + sa.String(length=3000, collation=get_asset_name_collation()), "mysql", ), nullable=False, diff --git a/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py b/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py index 39b8d8151692b..196176ed8537f 100644 --- a/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py +++ b/airflow-core/src/airflow/migrations/versions/0038_3_0_0_add_asset_active.py @@ -30,7 +30,7 @@ import sqlalchemy as sa from alembic import op -from airflow.migrations.utils import asset_name_collation +from airflow.migrations.utils import get_asset_name_collation # revision identifiers, used by Alembic. revision = "5a5d66100783" @@ -40,7 +40,7 @@ airflow_version = "3.0.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation=asset_name_collation()), + sa.String(length=1500, collation=get_asset_name_collation()), "mysql", ) diff --git a/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py b/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py index 79ded55d6d83d..4bee14bd1ffc0 100644 --- a/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py +++ b/airflow-core/src/airflow/migrations/versions/0039_3_0_0_tweak_assetaliasmodel_to_match_asset.py @@ -42,7 +42,7 @@ import sqlalchemy as sa from alembic import op -from airflow.migrations.utils import asset_name_collation +from airflow.migrations.utils import get_asset_name_collation # Revision identifiers, used by Alembic. revision = "fb2d4922cd79" @@ -52,7 +52,7 @@ airflow_version = "3.0.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation=asset_name_collation()), + sa.String(length=1500, collation=get_asset_name_collation()), "mysql", ) @@ -78,7 +78,7 @@ def downgrade(): batch_op.alter_column( "name", type_=sa.String(length=3000).with_variant( - sa.String(length=3000, collation=asset_name_collation()), + sa.String(length=3000, collation=get_asset_name_collation()), "mysql", ), nullable=False, diff --git a/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py b/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py index 0eb50af20adb7..b9a337a99ccb0 100644 --- a/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py +++ b/airflow-core/src/airflow/migrations/versions/0054_3_0_0_add_asset_reference_models.py @@ -30,7 +30,7 @@ from alembic import op from airflow.migrations.db_types import StringID -from airflow.migrations.utils import asset_name_collation +from airflow.migrations.utils import get_asset_name_collation from airflow.utils.sqlalchemy import UtcDateTime # revision identifiers, used by Alembic. @@ -41,7 +41,7 @@ airflow_version = "3.0.0" ASSET_STR_FIELD = sa.String(length=1500).with_variant( - sa.String(length=1500, collation=asset_name_collation()), "mysql" + sa.String(length=1500, collation=get_asset_name_collation()), "mysql" ) diff --git a/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py b/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py index 17ab84969d16a..73ef5bf6301ae 100644 --- a/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py +++ b/airflow-core/src/airflow/migrations/versions/0088_3_2_0_replace_asset_trigger_table_with_asset.py @@ -30,7 +30,7 @@ import sqlalchemy as sa from alembic import op -from airflow.migrations.utils import asset_name_collation +from airflow.migrations.utils import get_asset_name_collation # revision identifiers, used by Alembic. revision = "15d84ca19038" @@ -40,7 +40,7 @@ airflow_version = "3.2.0" _STRING_COLUMN_TYPE = sa.String(length=1500).with_variant( - sa.String(length=1500, collation=asset_name_collation()), + sa.String(length=1500, collation=get_asset_name_collation()), "mysql", ) diff --git a/airflow-core/src/airflow/models/base.py b/airflow-core/src/airflow/models/base.py index 2840bf7e8f8a8..3e219ed1284c1 100644 --- a/airflow-core/src/airflow/models/base.py +++ b/airflow-core/src/airflow/models/base.py @@ -84,15 +84,7 @@ def get_id_collation_args(): def get_asset_str_field(length: int = 1500) -> String: - """ - Build the string type used for asset name/uri/group columns. - - On MySQL these carry an explicit latin1 collation: the values are ASCII, and - a 1-byte-per-character charset keeps the 1500-char unique indexes inside the - 3072-byte index limit that utf8mb4 would blow past. The collation is - overridable because MySQL-compatible engines do not all ship - ``latin1_general_cs`` (TiDB, for one, accepts only ``latin1_bin``). - """ + """Build the string type used for indexed asset identifiers.""" collation = conf.get("database", "sql_engine_collation_for_asset_names", fallback="latin1_general_cs") return String(length=length).with_variant(String(length=length, collation=collation), "mysql") diff --git a/airflow-core/tests/unit/migrations/test_migration_utils.py b/airflow-core/tests/unit/migrations/test_migration_utils.py index 24a3ce580b455..f681cc9e54c04 100644 --- a/airflow-core/tests/unit/migrations/test_migration_utils.py +++ b/airflow-core/tests/unit/migrations/test_migration_utils.py @@ -46,6 +46,7 @@ from airflow import settings from airflow.migrations.utils import ( disable_sqlite_fkeys, + get_asset_name_collation, ignore_sqlite_value_error, mysql_drop_foreignkey_if_exists, ) @@ -56,6 +57,8 @@ upgradedb, ) +from tests_common.test_utils.config import conf_vars + # The stairway runs hundreds of upgrade/downgrade cycles. Each cycle goes # through ``_single_connection_pool`` which calls ``settings.reconfigure_orm()`` # at entry and exit, and ``dispose_orm()`` cannot synchronously close asyncpg @@ -67,6 +70,12 @@ pytestmark = pytest.mark.db_test + +def test_get_asset_name_collation(): + with conf_vars({("database", "sql_engine_collation_for_asset_names"): "latin1_bin"}): + assert get_asset_name_collation() == "latin1_bin" + + # Stairway starts from the 3.0.0 head revision. Starting here (rather than # the 2.6.2 squashed baseline) avoids triggering FAB-provider downgrade # handling that airflow.utils.db.downgrade() applies when the target is below diff --git a/airflow-core/tests/unit/models/test_base.py b/airflow-core/tests/unit/models/test_base.py index 51fd430dbab23..123eb24560bb6 100644 --- a/airflow-core/tests/unit/models/test_base.py +++ b/airflow-core/tests/unit/models/test_base.py @@ -17,8 +17,9 @@ from __future__ import annotations import pytest +from sqlalchemy.dialects import mysql, postgresql -from airflow.models.base import get_id_collation_args +from airflow.models.base import get_asset_str_field, get_id_collation_args from tests_common.test_utils.config import conf_vars @@ -50,3 +51,16 @@ def test_collation(dsn, expected, extra): with conf_vars({("database", "sql_alchemy_conn"): dsn, **extra}): assert expected == get_id_collation_args() + + +@pytest.mark.parametrize( + ("dialect", "collation", "expected"), + [ + pytest.param(mysql.dialect(), "latin1_general_cs", "latin1_general_cs", id="mysql-default"), + pytest.param(mysql.dialect(), "latin1_bin", "latin1_bin", id="mysql-override"), + pytest.param(postgresql.dialect(), "latin1_bin", None, id="postgres"), + ], +) +def test_asset_str_field_collation(dialect, collation, expected): + with conf_vars({("database", "sql_engine_collation_for_asset_names"): collation}): + assert get_asset_str_field().dialect_impl(dialect).collation == expected From c07ac97f3f019557798da40c4d09826cedd166ba Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Mon, 24 Aug 2026 01:59:43 -0400 Subject: [PATCH 4/5] Test asset collation in migration SQL Helper tests alone would not catch a migration reverting to the unsupported hard-coded collation. Render the full MySQL migration chain with the override and assert every emitted asset collation uses it. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../unit/migrations/test_migration_utils.py | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) diff --git a/airflow-core/tests/unit/migrations/test_migration_utils.py b/airflow-core/tests/unit/migrations/test_migration_utils.py index f681cc9e54c04..43bf365a0e28e 100644 --- a/airflow-core/tests/unit/migrations/test_migration_utils.py +++ b/airflow-core/tests/unit/migrations/test_migration_utils.py @@ -37,6 +37,7 @@ from functools import partial import pytest +from alembic import command from alembic.migration import MigrationContext from alembic.operations import Operations from alembic.script import ScriptDirectory @@ -76,6 +77,24 @@ def test_get_asset_name_collation(): assert get_asset_name_collation() == "latin1_bin" +def test_asset_name_collation_in_mysql_offline_migrations(capsys): + with conf_vars( + { + ("database", "sql_alchemy_conn"): "mysql+pymysql://root@localhost/airflow", + ("database", "sql_engine_collation_for_asset_names"): "latin1_bin", + } + ): + command.upgrade( + _get_alembic_config(), + f"base:{_REVISION_HEADS_MAP['3.4.0']}", + sql=True, + ) + + emitted_sql = capsys.readouterr().out + assert emitted_sql.count("latin1_bin") == 15 + assert "latin1_general_cs" not in emitted_sql + + # Stairway starts from the 3.0.0 head revision. Starting here (rather than # the 2.6.2 squashed baseline) avoids triggering FAB-provider downgrade # handling that airflow.utils.db.downgrade() applies when the target is below From e8e3527ab7312f558a88e39c570e076340a11021 Mon Sep 17 00:00:00 2001 From: 1fanwang <1fannnw@gmail.com> Date: Mon, 24 Aug 2026 02:16:30 -0400 Subject: [PATCH 5/5] Test asset collation in migration downgrades Two historical downgrades recreate indexed asset columns. Render both MySQL downgrade paths with the override so an unsupported hard-coded collation cannot return there. Signed-off-by: 1fanwang <1fannnw@gmail.com> --- .../unit/migrations/test_migration_utils.py | 29 +++++++++++++++++++ 1 file changed, 29 insertions(+) diff --git a/airflow-core/tests/unit/migrations/test_migration_utils.py b/airflow-core/tests/unit/migrations/test_migration_utils.py index 43bf365a0e28e..fd548ed0f621d 100644 --- a/airflow-core/tests/unit/migrations/test_migration_utils.py +++ b/airflow-core/tests/unit/migrations/test_migration_utils.py @@ -95,6 +95,35 @@ def test_asset_name_collation_in_mysql_offline_migrations(capsys): assert "latin1_general_cs" not in emitted_sql +@pytest.mark.parametrize( + ("current_revision", "target_revision"), + [ + pytest.param("fb2d4922cd79", "5a5d66100783", id="asset-alias-name"), + pytest.param("0d9e73a75ee4", "44eabb1904b4", id="asset-uri"), + ], +) +def test_asset_name_collation_in_mysql_offline_downgrades( + capsys, + current_revision, + target_revision, +): + with conf_vars( + { + ("database", "sql_alchemy_conn"): "mysql+pymysql://root@localhost/airflow", + ("database", "sql_engine_collation_for_asset_names"): "latin1_bin", + } + ): + command.downgrade( + _get_alembic_config(), + f"{current_revision}:{target_revision}", + sql=True, + ) + + emitted_sql = capsys.readouterr().out + assert emitted_sql.count("latin1_bin") == 1 + assert "latin1_general_cs" not in emitted_sql + + # Stairway starts from the 3.0.0 head revision. Starting here (rather than # the 2.6.2 squashed baseline) avoids triggering FAB-provider downgrade # handling that airflow.utils.db.downgrade() applies when the target is below