kaxil commented on code in PR #74275:
URL: https://github.com/apache/airflow/pull/74275#discussion_r4222208774
##########
airflow-core/docs/core-concepts/task-state-store.rst:
##########
@@ -30,7 +30,7 @@ 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.
Review Comment:
`region_id` is now listed as part of the task scope, but none of the docs
say what a region is or when it differs from the all-zero default. One sentence
here would help, since someone with a custom backend reading this has to decide
whether their existing refs are still unique.
##########
airflow-core/src/airflow/utils/sqlalchemy.py:
##########
@@ -291,6 +292,58 @@ def load_dialect_impl(self, dialect):
return super().load_dialect_impl(dialect)
+_MYSQL_DIALECT_NAMES = ("mysql", "mariadb")
+
+
+class CompactUUID(TypeDecorator):
+ """
+ A UUID that MySQL stores as ``BINARY(16)``; other dialects use
SQLAlchemy's :class:`~sqlalchemy.types.Uuid`.
+
+ The default MySQL form, ``CHAR(32)``, takes 128 bytes in an index key
under a utf8mb4 collation.
+ Keys that already approach InnoDB's 3072-byte limit, such as those holding
several ``StringID`` columns,
+ overflow with it.
+ """
+
+ impl = Uuid
+ cache_ok = True
+
+ def load_dialect_impl(self, dialect):
+ if dialect.name in _MYSQL_DIALECT_NAMES:
+ return dialect.type_descriptor(mysql.BINARY(16))
+ return dialect.type_descriptor(Uuid())
+
+ def process_bind_param(self, value, dialect):
+ if value is None or dialect.name not in _MYSQL_DIALECT_NAMES:
+ return value
+ return (UUID(value) if isinstance(value, str) else value).bytes
+
+ def process_result_value(self, value, dialect):
+ if value is None or dialect.name not in _MYSQL_DIALECT_NAMES:
+ return value
+ return UUID(bytes=bytes(value))
+
+
+class compact_uuid_default(ColumnElement):
+ """Server default for a :class:`CompactUUID` column, rendered in the form
each dialect stores."""
+
+ inherit_cache = True
+ type = NullType()
+
+ def __init__(self, value: UUID):
+ self.value = value
+
+
+@compiles(compact_uuid_default)
+def _compact_uuid_default(element, compiler, **kw):
+ return compiler.render_literal_value(element.value.hex, String())
+
+
+@compiles(compact_uuid_default, "mysql")
+@compiles(compact_uuid_default, "mariadb")
+def _compact_uuid_default_mysql(element, compiler, **kw):
+ return f"UNHEX({compiler.render_literal_value(element.value.hex,
String())})"
Review Comment:
In offline mode (`airflow db migrate --show-sql-only`) SQLAlchemy has no
server version, so it doesn't add the parentheses MySQL needs around an
expression default. Compiling an `ADD COLUMN` with this default against a bare
`mysql.dialect()` gives `NOT NULL DEFAULT UNHEX('000...0')`; only once
`server_version_info` is set does it become `DEFAULT (UNHEX(...))`. MySQL
rejects the unparenthesized form, so the generated script stops at the first
`region_id` column after `dynamic_region` has already been created. A hex
literal (`X'000...0'`) needs no parentheses, and
`test_server_default_per_dialect` only compiles the bare expression, so a test
that compiles the column DDL against a MySQL dialect with no server version
would pin this.
##########
airflow-core/src/airflow/migrations/versions/0143_3_4_0_add_dynamic_region_storage.py:
##########
@@ -0,0 +1,215 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements. See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership. The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied. See the License for the
+# specific language governing permissions and limitations
+# under the License.
+
+"""
+Add dynamic region storage.
+
+Revision ID: 54a27b6f9d01
+Revises: e7c2a91bd540
+Create Date: 2026-09-27 20:00:00.000000
+"""
+
+from __future__ import annotations
+
+from textwrap import dedent
+from uuid import UUID
+
+import sqlalchemy as sa
+from alembic import op
+
+from airflow.migrations.utils import raise_if_rows_exist
+from airflow.models.base import StringID
+from airflow.utils.sqlalchemy import CompactUUID, UtcDateTime,
compact_uuid_default
+
+revision = "54a27b6f9d01"
+down_revision = "e7c2a91bd540"
+branch_labels = None
+depends_on = None
+airflow_version = "3.4.0"
+
+_SENTINEL = UUID(int=0)
+_KEYS = (
+ (
+ "task_instance",
+ "task_instance_current_key",
+ ("dag_id", "task_id", "run_id", "map_index", "working_set"),
+ "working_set IS NOT NULL",
+ ),
+ (
+ "task_instance",
+ "task_instance_try_key",
+ ("dag_id", "task_id", "run_id", "map_index", "try_number"),
+ None,
+ ),
+ ("task_state_store", "task_state_store_uq", ("dag_run_id", "task_id",
"map_index", "key"), None),
+)
+
+
+def _replace_unique(table_name, constraint_name, columns):
+ context = op.get_context()
+ dialect = context.dialect.name
+ if dialect in ("postgresql", "mysql"):
+ quote = context.dialect.identifier_preparer.quote
+ column_list = ", ".join(quote(name) for name in columns)
+ if dialect == "postgresql":
+ sql = f"""
+ ALTER TABLE {table_name}
+ DROP CONSTRAINT {constraint_name},
+ ADD CONSTRAINT {constraint_name} UNIQUE ({column_list})
+ """
+ else:
+ sql = f"""
+ ALTER TABLE {table_name}
+ DROP INDEX {constraint_name},
+ ADD UNIQUE KEY {constraint_name} ({column_list})
+ """
+ op.execute(dedent(sql))
+ else:
+ with op.get_context().autocommit_block():
Review Comment:
On SQLite this `autocommit_block()` commits everything the migration has
done so far (`dynamic_region`, its indexes and both `region_id` columns), and
the batch rebuild then runs outside a transaction. If a later rebuild fails,
the DB is left half-migrated and a rerun stops on `dynamic_region` already
existing. `_sqlite_rebuilds()` in 0142 keeps only the PRAGMAs in autocommit and
runs the rebuild under `begin_nested()`; could that move to
`migrations/utils.py` and be reused here?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]