diff --git a/airflow/api/common/trigger_dag.py b/airflow/api/common/trigger_dag.py index b18957261f3a0..57866ae09c8c0 100644 --- a/airflow/api/common/trigger_dag.py +++ b/airflow/api/common/trigger_dag.py @@ -58,6 +58,8 @@ def _trigger_dag( :param replace_microseconds: whether microseconds should be zeroed :return: list of triggered dags """ + from airflow.models.serialized_dag import SerializedDagModel + dag = dag_bag.get_dag(dag_id) # prefetch dag if it is stored serialized if dag is None or dag_id not in dag_bag.dags: @@ -99,7 +101,7 @@ def _trigger_dag( state=DagRunState.QUEUED, conf=run_conf, external_trigger=True, - dag_hash=dag_bag.dags_hash.get(dag_id), + serialized_dag=SerializedDagModel.get(dag_id), data_interval=data_interval, triggered_by=triggered_by, ) diff --git a/airflow/api_connexion/endpoints/dag_run_endpoint.py b/airflow/api_connexion/endpoints/dag_run_endpoint.py index 44891c0ef2c84..490e08800c50e 100644 --- a/airflow/api_connexion/endpoints/dag_run_endpoint.py +++ b/airflow/api_connexion/endpoints/dag_run_endpoint.py @@ -61,6 +61,7 @@ from airflow.auth.managers.models.resource_details import DagAccessEntity from airflow.exceptions import ParamValidationError from airflow.models import DagModel, DagRun +from airflow.models.serialized_dag import SerializedDagModel from airflow.timetables.base import DataInterval from airflow.utils.airflow_flask_app import get_airflow_app from airflow.utils.db import get_query_count @@ -347,7 +348,7 @@ def post_dag_run(*, dag_id: str, session: Session = NEW_SESSION) -> APIResponse: state=DagRunState.QUEUED, conf=post_body.get("conf"), external_trigger=True, - dag_hash=get_airflow_app().dag_bag.dags_hash.get(dag_id), + serialized_dag=SerializedDagModel.get(dag_id), session=session, triggered_by=DagRunTriggeredByType.REST_API, ) diff --git a/airflow/jobs/scheduler_job_runner.py b/airflow/jobs/scheduler_job_runner.py index 21c871e951253..2e6b84c722c09 100644 --- a/airflow/jobs/scheduler_job_runner.py +++ b/airflow/jobs/scheduler_job_runner.py @@ -1318,7 +1318,7 @@ def _create_dag_runs(self, dag_models: Collection[DagModel], session: Session) - self.log.error("DAG '%s' not found in serialized_dag table", dag_model.dag_id) continue - dag_hash = self.dagbag.dags_hash.get(dag.dag_id) + serialized_dag = SerializedDagModel.get(dag.dag_id, session=session) data_interval = dag.get_next_data_interval(dag_model) # Explicitly check if the DagRun already exists. This is an edge case @@ -1338,7 +1338,7 @@ def _create_dag_runs(self, dag_models: Collection[DagModel], session: Session) - data_interval=data_interval, external_trigger=False, session=session, - dag_hash=dag_hash, + serialized_dag=serialized_dag, creating_job_id=self.job.id, triggered_by=DagRunTriggeredByType.TIMETABLE, ) @@ -1397,7 +1397,7 @@ def _create_dag_runs_asset_triggered( ) continue - dag_hash = self.dagbag.dags_hash.get(dag.dag_id) + serialized_dag = SerializedDagModel.get(dag.dag_id, session=session) # Explicitly check if the DagRun already exists. This is an edge case # where a Dag Run is created but `DagModel.next_dagrun` and `DagModel.next_dagrun_create_after` @@ -1452,7 +1452,7 @@ def _create_dag_runs_asset_triggered( state=DagRunState.QUEUED, external_trigger=False, session=session, - dag_hash=dag_hash, + serialized_dag=serialized_dag, creating_job_id=self.job.id, triggered_by=DagRunTriggeredByType.DATASET, ) @@ -1701,12 +1701,13 @@ def _verify_integrity_if_dag_changed(self, dag_run: DagRun, session: Session) -> Return True if we determine that DAG still exists. """ - latest_version = SerializedDagModel.get_latest_version_hash(dag_run.dag_id, session=session) - if dag_run.dag_hash == latest_version: + latest_version = SerializedDagModel.get(dag_run.dag_id, session=session) + + if latest_version and dag_run.serialized_dag_id == latest_version.id: self.log.debug("DAG %s not changed structure, skipping dagrun.verify_integrity", dag_run.dag_id) return True - dag_run.dag_hash = latest_version + dag_run.serialized_dag = latest_version # Refresh the DAG dag_run.dag = self.dagbag.get_dag(dag_id=dag_run.dag_id, session=session) diff --git a/airflow/migrations/versions/0036_3_0_0_add_serial_id_to_sdm.py b/airflow/migrations/versions/0036_3_0_0_add_serial_id_to_sdm.py new file mode 100644 index 0000000000000..5ce5e39249082 --- /dev/null +++ b/airflow/migrations/versions/0036_3_0_0_add_serial_id_to_sdm.py @@ -0,0 +1,66 @@ +# +# 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 serial ID to SDM. + +Revision ID: e1ff90d3efe9 +Revises: 0d9e73a75ee4 +Create Date: 2024-09-27 09:32:46.514067 + +""" + +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op + +from airflow.models.base import naming_convention + +# revision identifiers, used by Alembic. +revision = "e1ff90d3efe9" +down_revision = "0d9e73a75ee4" +branch_labels = None +depends_on = None +airflow_version = "3.0.0" + + +def upgrade(): + """Apply add serial pkey to SerializedDag.""" + with op.batch_alter_table( + "serialized_dag", recreate="always", naming_convention=naming_convention + ) as batch_op: + batch_op.drop_constraint("serialized_dag_pkey", type_="primary") + # hack. The primary_key here sets autoincrement + batch_op.add_column(sa.Column("id", sa.Integer(), primary_key=True), insert_before="dag_id") + batch_op.create_primary_key("serialized_dag_pkey", ["id"]) + batch_op.add_column( + sa.Column("version_number", sa.Integer(), nullable=False, default=1), insert_before="dag_id" + ) + batch_op.create_unique_constraint( + batch_op.f("dag_hash_version_number_unique"), ["dag_hash", "version_number"] + ) + + +def downgrade(): + """Unapply add serial pkey to SerializedDag.""" + with op.batch_alter_table("serialized_dag", naming_convention=naming_convention) as batch_op: + batch_op.drop_constraint(batch_op.f("dag_hash_version_number_unique"), type_="unique") + batch_op.drop_column("id") + batch_op.create_primary_key("serialized_dag_pkey", ["dag_id"]) + batch_op.drop_column("version_number") diff --git a/airflow/migrations/versions/0037_3_0_0_add_sdm_foreignkey_to_dagrun.py b/airflow/migrations/versions/0037_3_0_0_add_sdm_foreignkey_to_dagrun.py new file mode 100644 index 0000000000000..8ae5c714bd609 --- /dev/null +++ b/airflow/migrations/versions/0037_3_0_0_add_sdm_foreignkey_to_dagrun.py @@ -0,0 +1,80 @@ +# +# 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 SDM foreign key to DagRun, TI & TIH. + +Revision ID: 4235395d5ec5 +Revises: e1ff90d3efe9 +Create Date: 2024-10-03 13:37:55.678831 + +""" + +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op + +# revision identifiers, used by Alembic. +revision = "4235395d5ec5" +down_revision = "e1ff90d3efe9" +branch_labels = None +depends_on = None +airflow_version = "3.0.0" + + +def upgrade(): + """Apply Add SDM foreignkey to DagRun.""" + with op.batch_alter_table("dag_run") as batch_op: + batch_op.add_column(sa.Column("serialized_dag_id", sa.Integer())) + batch_op.create_foreign_key( + "dag_run_serialized_dag_fkey", + "serialized_dag", + ["serialized_dag_id"], + ["id"], + ondelete="SET NULL", + ) + batch_op.drop_column("dag_hash") + + with op.batch_alter_table("task_instance") as batch_op: + batch_op.add_column(sa.Column("serialized_dag_id", sa.Integer())) + batch_op.create_foreign_key( + "task_instance_serialized_dag_fkey", + "serialized_dag", + ["serialized_dag_id"], + ["id"], + ondelete="SET NULL", + ) + + with op.batch_alter_table("task_instance_history") as batch_op: + batch_op.add_column(sa.Column("serialized_dag_id", sa.Integer())) + + +def downgrade(): + """Unapply Add SDM foreignkey to DagRun.""" + with op.batch_alter_table("dag_run") as batch_op: + batch_op.add_column(sa.Column("dag_hash", sa.String(32))) + batch_op.drop_constraint("dag_run_serialized_dag_fkey", type_="foreignkey") + batch_op.drop_column("serialized_dag_id") + + with op.batch_alter_table("task_instance") as batch_op: + batch_op.drop_constraint("task_instance_serialized_dag_fkey", type_="foreignkey") + batch_op.drop_column("serialized_dag_id") + + with op.batch_alter_table("task_instance_history") as batch_op: + batch_op.drop_column("serialized_dag_id") diff --git a/airflow/models/backfill.py b/airflow/models/backfill.py index db10c804aac0d..5514a6ce6db46 100644 --- a/airflow/models/backfill.py +++ b/airflow/models/backfill.py @@ -124,7 +124,7 @@ def _create_backfill( dag_run_conf: dict | None, ) -> Backfill | None: with create_session() as session: - serdag = session.get(SerializedDagModel, dag_id) + serdag = session.scalar(SerializedDagModel.latest_item_select_object(dag_id)) if not serdag: raise NotFound(f"Could not find dag {dag_id}") diff --git a/airflow/models/dag.py b/airflow/models/dag.py index 0632819952ae4..49cdb33dfe9ab 100644 --- a/airflow/models/dag.py +++ b/airflow/models/dag.py @@ -302,7 +302,7 @@ def _create_orm_dagrun( conf, state, run_type, - dag_hash, + serialized_dag, creating_job_id, data_interval, session, @@ -317,7 +317,7 @@ def _create_orm_dagrun( conf=conf, state=state, run_type=run_type, - dag_hash=dag_hash, + serialized_dag=serialized_dag, creating_job_id=creating_job_id, data_interval=data_interval, triggered_by=triggered_by, @@ -2542,7 +2542,7 @@ def create_dagrun( conf: dict | None = None, run_type: DagRunType | None = None, session: Session = NEW_SESSION, - dag_hash: str | None = None, + serialized_dag: SerializedDagModel | None = None, creating_job_id: int | None = None, data_interval: tuple[datetime, datetime] | None = None, ): @@ -2561,7 +2561,7 @@ def create_dagrun( :param conf: Dict containing configuration/parameters to pass to the DAG :param creating_job_id: id of the job creating this DagRun :param session: database session - :param dag_hash: Hash of Serialized DAG + :param serialized_dag: The serialized Dag Model object :param data_interval: Data interval of the DagRun """ logical_date = timezone.coerce_datetime(execution_date) @@ -2627,7 +2627,7 @@ def create_dagrun( conf=conf, state=state, run_type=run_type, - dag_hash=dag_hash, + serialized_dag=serialized_dag, creating_job_id=creating_job_id, data_interval=data_interval, session=session, diff --git a/airflow/models/dagrun.py b/airflow/models/dagrun.py index d1dbeaf41e58b..20ecd170db7c9 100644 --- a/airflow/models/dagrun.py +++ b/airflow/models/dagrun.py @@ -80,6 +80,7 @@ from airflow.models.dag import DAG from airflow.models.operator import Operator + from airflow.models.serialized_dag import SerializedDagModel from airflow.serialization.pydantic.dag_run import DagRunPydantic from airflow.serialization.pydantic.taskinstance import TaskInstancePydantic from airflow.serialization.pydantic.tasklog import LogTemplatePydantic @@ -141,7 +142,6 @@ class DagRun(Base, LoggingMixin): data_interval_end = Column(UtcDateTime) # When a scheduler last attempted to schedule TIs for this DagRun last_scheduling_decision = Column(UtcDateTime) - dag_hash = Column(String(32)) # Foreign key to LogTemplate. DagRun rows created prior to this column's # existence have this set to NULL. Later rows automatically populate this on # insert to point to the latest LogTemplate entry. @@ -155,6 +155,11 @@ class DagRun(Base, LoggingMixin): # This number is incremented only when the DagRun is re-Queued, # when the DagRun is cleared. clear_number = Column(Integer, default=0, nullable=False, server_default="0") + serialized_dag_id = Column( + Integer, + ForeignKey("serialized_dag.id", name="dag_run_serialized_dag_fkey", ondelete="SET NULL"), + ) + serialized_dag = relationship("SerializedDagModel", back_populates="dag_run") # Remove this `if` after upgrading Sphinx-AutoAPI if not TYPE_CHECKING and "BUILDING_AIRFLOW_DOCS" in os.environ: @@ -218,7 +223,7 @@ def __init__( conf: Any | None = None, state: DagRunState | None = None, run_type: str | None = None, - dag_hash: str | None = None, + serialized_dag: SerializedDagModel | None = None, creating_job_id: int | None = None, data_interval: tuple[datetime, datetime] | None = None, triggered_by: DagRunTriggeredByType | None = None, @@ -242,7 +247,7 @@ def __init__( else: self.queued_at = queued_at self.run_type = run_type - self.dag_hash = dag_hash + self.serialized_dag = serialized_dag self.creating_job_id = creating_job_id self.clear_number = 0 self.triggered_by = triggered_by @@ -354,6 +359,12 @@ def set_state(self, state: DagRunState) -> None: def state(self): return synonym("_state", descriptor=property(self.get_state, self.set_state)) + @provide_session + def dag_hash(self, session: Session = NEW_SESSION): + from airflow.models.serialized_dag import SerializedDagModel as SDM + + return str(session.scalar(select(SDM.dag_hash).where(SDM.id == self.serialized_dag_id))) + @provide_session def refresh_from_db(self, session: Session = NEW_SESSION) -> None: """ @@ -939,6 +950,7 @@ def recalculate(self) -> _UnfinishedStates: "state=%s, external_trigger=%s, run_type=%s, " "data_interval_start=%s, data_interval_end=%s, dag_hash=%s" ) + self.log.info( msg, self.dag_id, @@ -956,7 +968,7 @@ def recalculate(self) -> _UnfinishedStates: self.run_type, self.data_interval_start, self.data_interval_end, - self.dag_hash, + self.dag_hash(session), ) with Trace.start_span_from_dagrun(dagrun=self) as span: @@ -980,7 +992,7 @@ def recalculate(self) -> _UnfinishedStates: "run_type": str(self.run_type), "data_interval_start": str(self.data_interval_start), "data_interval_end": str(self.data_interval_end), - "dag_hash": str(self.dag_hash), + "dag_hash": str(self.dag_hash(session)), "conf": str(self.conf), } if span.is_recording(): diff --git a/airflow/models/serialized_dag.py b/airflow/models/serialized_dag.py index 32be31d721e34..55a710ae1e3d9 100644 --- a/airflow/models/serialized_dag.py +++ b/airflow/models/serialized_dag.py @@ -25,8 +25,21 @@ from typing import TYPE_CHECKING, Any, Collection import sqlalchemy_jsonfield -from sqlalchemy import BigInteger, Column, Index, LargeBinary, String, and_, exc, or_, select -from sqlalchemy.orm import backref, foreign, relationship +from sqlalchemy import ( + BigInteger, + Column, + Index, + Integer, + LargeBinary, + String, + UniqueConstraint, + and_, + delete, + exc, + or_, + select, +) +from sqlalchemy.orm import backref, relationship from sqlalchemy.sql.expression import func, literal from airflow.api_internal.internal_api_call import internal_api_call @@ -34,7 +47,6 @@ from airflow.models.base import ID_LEN, Base from airflow.models.dag import DagModel from airflow.models.dagcode import DagCode -from airflow.models.dagrun import DagRun from airflow.serialization.dag_dependency import DagDependency from airflow.serialization.serialized_objects import SerializedDAG from airflow.settings import COMPRESS_SERIALIZED_DAGS, MIN_SERIALIZED_DAG_UPDATE_INTERVAL, json @@ -77,7 +89,9 @@ class SerializedDagModel(Base): __tablename__ = "serialized_dag" - dag_id = Column(String(ID_LEN), primary_key=True) + id = Column(Integer, primary_key=True) + version_number = Column(Integer, nullable=False, default=1) + dag_id = Column(String(ID_LEN)) fileloc = Column(String(2000), nullable=False) # The max length of fileloc exceeds the limit of indexing. fileloc_hash = Column(BigInteger(), nullable=False) @@ -87,14 +101,14 @@ class SerializedDagModel(Base): dag_hash = Column(String(32), nullable=False) processor_subdir = Column(String(2000), nullable=True) - __table_args__ = (Index("idx_fileloc_hash", fileloc_hash, unique=False),) - - dag_runs = relationship( - DagRun, - primaryjoin=dag_id == foreign(DagRun.dag_id), # type: ignore - backref=backref("serialized_dag", uselist=False, innerjoin=True), + __table_args__ = ( + Index("idx_fileloc_hash", fileloc_hash, unique=False), + UniqueConstraint("dag_hash", "version_number", name="dag_hash_version_number_unique"), ) + dag_run = relationship("DagRun", back_populates="serialized_dag") + task_instance = relationship("TaskInstance", back_populates="serialized_dag") + dag_model = relationship( DagModel, primaryjoin=dag_id == DagModel.dag_id, # type: ignore @@ -193,8 +207,9 @@ def write_dag( log.debug("Checking if DAG (%s) changed", dag.dag_id) new_serialized_dag = cls(dag, processor_subdir) + serialized_dag_db = session.execute( - select(cls.dag_hash, cls.processor_subdir).where(cls.dag_id == dag.dag_id) + select(cls.dag_hash, cls.processor_subdir).where(cls.dag_id == dag.dag_id).order_by(cls.id.desc()) ).first() if ( @@ -205,11 +220,24 @@ def write_dag( log.debug("Serialized DAG (%s) is unchanged. Skipping writing to DB", dag.dag_id) return False + latest_version = session.execute( + select(cls.version_number).where(cls.dag_id == dag.dag_id).order_by(cls.id.desc()).limit(1) + ).first() + if latest_version is not None: + new_version_number = latest_version.version_number + 1 + else: + new_version_number = 1 + new_serialized_dag.version_number = new_version_number + log.debug("Writing Serialized DAG: %s to the DB", dag.dag_id) - session.merge(new_serialized_dag) + session.add(new_serialized_dag) log.debug("DAG: %s written to the DB", dag.dag_id) return True + @classmethod + def latest_item_select_object(cls, dag_id): + return select(cls).where(cls.dag_id == dag_id).order_by(cls.id.desc()).limit(1) + @classmethod @provide_session def read_all_dags(cls, session: Session = NEW_SESSION) -> dict[str, SerializedDAG]: @@ -219,7 +247,18 @@ def read_all_dags(cls, session: Session = NEW_SESSION) -> dict[str, SerializedDA :param session: ORM Session :returns: a dict of DAGs read from database """ - serialized_dags = session.scalars(select(cls)) + latest_versions_subquery = ( + session.query(cls.dag_id, func.max(cls.version_number).label("max_version")) + .group_by(cls.dag_id) + .subquery() + ) + serialized_dags = session.scalars( + select(cls).join( + latest_versions_subquery, + (cls.dag_id == latest_versions_subquery.c.dag_id) + and (cls.version_number == latest_versions_subquery.c.max_version), + ) + ) dags = {} for row in serialized_dags: @@ -269,7 +308,7 @@ def remove_dag(cls, dag_id: str, session: Session = NEW_SESSION) -> None: :param dag_id: dag_id to be deleted :param session: ORM Session. """ - session.execute(cls.__table__.delete().where(cls.dag_id == dag_id)) + session.execute(delete(cls).where(cls.dag_id == dag_id)) @classmethod @internal_api_call @@ -334,11 +373,7 @@ def get(cls, dag_id: str, session: Session = NEW_SESSION) -> SerializedDagModel :param dag_id: the DAG to fetch :param session: ORM Session """ - row = session.scalar(select(cls).where(cls.dag_id == dag_id)) - if row: - return row - - return session.scalar(select(cls).where(cls.dag_id == dag_id)) + return session.scalar(cls.latest_item_select_object(dag_id)) @staticmethod @provide_session @@ -373,7 +408,9 @@ def get_last_updated_datetime(cls, dag_id: str, session: Session = NEW_SESSION) :param dag_id: DAG ID :param session: ORM Session """ - return session.scalar(select(cls.last_updated).where(cls.dag_id == dag_id)) + return session.scalar( + select(cls.last_updated).where(cls.dag_id == dag_id).order_by(cls.id.desc()).limit(1) + ) @classmethod @provide_session @@ -395,7 +432,9 @@ def get_latest_version_hash(cls, dag_id: str, session: Session = NEW_SESSION) -> :param session: ORM Session :return: DAG Hash, or None if the DAG is not found """ - return session.scalar(select(cls.dag_hash).where(cls.dag_id == dag_id)) + return session.scalar( + select(cls.dag_hash).where(cls.dag_id == dag_id).order_by(cls.id.desc()).limit(1) + ) @classmethod def get_latest_version_hash_and_updated_datetime( @@ -413,7 +452,10 @@ def get_latest_version_hash_and_updated_datetime( :return: A tuple of DAG Hash and last updated datetime, or None if the DAG is not found """ return session.execute( - select(cls.dag_hash, cls.last_updated).where(cls.dag_id == dag_id) + select(cls.dag_hash, cls.last_updated) + .where(cls.dag_id == dag_id) + .order_by(cls.id.desc()) + .limit(1) ).one_or_none() @classmethod @@ -424,9 +466,18 @@ def get_dag_dependencies(cls, session: Session = NEW_SESSION) -> dict[str, list[ :param session: ORM Session """ + latest_versions_subquery = ( + session.query(cls.dag_id, func.max(cls.version_number).label("max_version")) + .group_by(cls.dag_id) + .subquery() + ) if session.bind.dialect.name in ["sqlite", "mysql"]: query = session.execute( - select(cls.dag_id, func.json_extract(cls._data, "$.dag.dag_dependencies")) + select(cls.dag_id, func.json_extract(cls._data, "$.dag.dag_dependencies")).join( + latest_versions_subquery, + (cls.dag_id == latest_versions_subquery.c.dag_id) + and (cls.version_number == latest_versions_subquery.c.max_version), + ) ) iterator = ((dag_id, json.loads(deps_data) if deps_data else []) for dag_id, deps_data in query) else: @@ -439,10 +490,9 @@ def get_dag_dependencies(cls, session: Session = NEW_SESSION) -> dict[str, list[ @internal_api_call @provide_session def get_serialized_dag(dag_id: str, task_id: str, session: Session = NEW_SESSION) -> Operator | None: - from airflow.models.serialized_dag import SerializedDagModel - try: - model = session.get(SerializedDagModel, dag_id) + # get the latest version of the DAG + model = session.scalar(SerializedDagModel.latest_item_select_object(dag_id)) if model: return model.dag.get_task(task_id) except (exc.NoResultFound, TaskNotFound): diff --git a/airflow/models/taskinstance.py b/airflow/models/taskinstance.py index 333a4cad91cbe..be18f31b11d2a 100644 --- a/airflow/models/taskinstance.py +++ b/airflow/models/taskinstance.py @@ -43,6 +43,7 @@ Column, DateTime, Float, + ForeignKey, ForeignKeyConstraint, Index, Integer, @@ -1858,8 +1859,11 @@ class TaskInstance(Base, LoggingMixin): next_kwargs = Column(MutableDict.as_mutable(ExtendedJSON)) _task_display_property_value = Column("task_display_name", String(2000), nullable=True) - # If adding new fields here then remember to add them to - # refresh_from_db() or they won't display in the UI correctly + serialized_dag_id = Column( + Integer, + ForeignKey("serialized_dag.id", name="task_instance_serialized_dag_fkey", ondelete="SET NULL"), + ) + serialized_dag = relationship("SerializedDagModel", back_populates="task_instance") __table_args__ = ( Index("ti_dag_state", dag_id, state), diff --git a/airflow/models/taskinstancehistory.py b/airflow/models/taskinstancehistory.py index ccdca700af6e9..4f15ac4747ee4 100644 --- a/airflow/models/taskinstancehistory.py +++ b/airflow/models/taskinstancehistory.py @@ -92,6 +92,7 @@ class TaskInstanceHistory(Base): next_kwargs = Column(MutableDict.as_mutable(ExtendedJSON)) task_display_name = Column("task_display_name", String(2000), nullable=True) + serialized_dag_id = Column(Integer) def __init__( self, @@ -100,10 +101,9 @@ def __init__( ): super().__init__() for column in self.__table__.columns: - if column.name == "id": + if column.name in ["id", "dag_hash"]: continue setattr(self, column.name, getattr(ti, column.name)) - if state: self.state = state diff --git a/airflow/serialization/pydantic/dag_run.py b/airflow/serialization/pydantic/dag_run.py index 86857452e8310..1f3abdad53e75 100644 --- a/airflow/serialization/pydantic/dag_run.py +++ b/airflow/serialization/pydantic/dag_run.py @@ -52,7 +52,7 @@ class DagRunPydantic(BaseModelPydantic): data_interval_start: Optional[datetime] data_interval_end: Optional[datetime] last_scheduling_decision: Optional[datetime] - dag_hash: Optional[str] + serialized_dag_id: Optional[int] updated_at: Optional[datetime] dag: Optional[PydanticDag] consumed_dataset_events: List[AssetEventPydantic] # noqa: UP006 diff --git a/airflow/utils/db.py b/airflow/utils/db.py index 898da02078d5a..8bccb3942e639 100644 --- a/airflow/utils/db.py +++ b/airflow/utils/db.py @@ -96,7 +96,7 @@ class MappedClassProtocol(Protocol): "2.9.0": "1949afb29106", "2.9.2": "686269002441", "2.10.0": "22ed7efa9da2", - "3.0.0": "0d9e73a75ee4", + "3.0.0": "4235395d5ec5", } diff --git a/airflow/www/views.py b/airflow/www/views.py index 47c548d5e7667..84838985200b1 100644 --- a/airflow/www/views.py +++ b/airflow/www/views.py @@ -2254,7 +2254,7 @@ def trigger(self, dag_id: str, session: Session = NEW_SESSION): state=DagRunState.QUEUED, conf=run_conf, external_trigger=True, - dag_hash=get_airflow_app().dag_bag.dags_hash.get(dag_id), + serialized_dag=SerializedDagModel.get(dag.dag_id), run_id=run_id, triggered_by=DagRunTriggeredByType.UI, ) diff --git a/docs/apache-airflow/img/airflow_erd.sha256 b/docs/apache-airflow/img/airflow_erd.sha256 index bca068fde6749..48830b1ddaff8 100644 --- a/docs/apache-airflow/img/airflow_erd.sha256 +++ b/docs/apache-airflow/img/airflow_erd.sha256 @@ -1 +1 @@ -64dfad12dfd49f033c4723c2f3bb3bac58dd956136fb24a87a2e5a6ae176ec1a \ No newline at end of file +3fe5f6994c5ccdada9ce9834d3ee314cef70c243f968aba4e399e6e1f416d720 \ No newline at end of file diff --git a/docs/apache-airflow/img/airflow_erd.svg b/docs/apache-airflow/img/airflow_erd.svg index 4eb6c2ee70917..a43b76f42e987 100644 --- a/docs/apache-airflow/img/airflow_erd.svg +++ b/docs/apache-airflow/img/airflow_erd.svg @@ -4,11 +4,11 @@ - - + + %3 - + log @@ -360,2015 +360,2060 @@ [TEXT] - - -backfill - -backfill - -id - - [INTEGER] - NOT NULL - -completed_at - - [TIMESTAMP] - -created_at - - [TIMESTAMP] - NOT NULL - -dag_id - - [VARCHAR(250)] - NOT NULL - -dag_run_conf - - [JSON] - -from_date - - [TIMESTAMP] - NOT NULL - -is_paused - - [BOOLEAN] - -max_active_runs - - [INTEGER] - NOT NULL - -to_date - - [TIMESTAMP] - NOT NULL - -updated_at - - [TIMESTAMP] - NOT NULL - - - -backfill_dag_run - -backfill_dag_run - -id - - [INTEGER] - NOT NULL - -backfill_id - - [INTEGER] - NOT NULL - -dag_run_id - - [INTEGER] - -sort_ordinal - - [INTEGER] - NOT NULL - - + import_error - -import_error - -id - - [INTEGER] - NOT NULL - -filename - - [VARCHAR(1024)] - -processor_subdir - - [VARCHAR(2000)] - -stacktrace - - [TEXT] - -timestamp - - [TIMESTAMP] - - - -serialized_dag - -serialized_dag - -dag_id - - [VARCHAR(250)] - NOT NULL - -dag_hash - - [VARCHAR(32)] - NOT NULL - -data - - [JSON] - -data_compressed - - [BYTEA] - -fileloc - - [VARCHAR(2000)] - NOT NULL - -fileloc_hash - - [BIGINT] - NOT NULL - -last_updated - - [TIMESTAMP] - NOT NULL - -processor_subdir - - [VARCHAR(2000)] + +import_error + +id + + [INTEGER] + NOT NULL + +filename + + [VARCHAR(1024)] + +processor_subdir + + [VARCHAR(2000)] + +stacktrace + + [TEXT] + +timestamp + + [TIMESTAMP] - + dataset_alias - -dataset_alias - -id - - [INTEGER] - NOT NULL - -name - - [VARCHAR(3000)] - NOT NULL + +dataset_alias + +id + + [INTEGER] + NOT NULL + +name + + [VARCHAR(3000)] + NOT NULL - + dataset_alias_dataset - -dataset_alias_dataset - -alias_id - - [INTEGER] - NOT NULL - -dataset_id - - [INTEGER] - NOT NULL + +dataset_alias_dataset + +alias_id + + [INTEGER] + NOT NULL + +dataset_id + + [INTEGER] + NOT NULL dataset_alias--dataset_alias_dataset - -0..N -1 + +0..N +1 dataset_alias--dataset_alias_dataset - -0..N -1 + +0..N +1 - + dataset_alias_dataset_event - -dataset_alias_dataset_event - -alias_id - - [INTEGER] - NOT NULL - -event_id - - [INTEGER] - NOT NULL + +dataset_alias_dataset_event + +alias_id + + [INTEGER] + NOT NULL + +event_id + + [INTEGER] + NOT NULL dataset_alias--dataset_alias_dataset_event - -0..N -1 + +0..N +1 dataset_alias--dataset_alias_dataset_event - -0..N -1 + +0..N +1 - + dag_schedule_dataset_alias_reference - -dag_schedule_dataset_alias_reference - -alias_id - - [INTEGER] - NOT NULL - -dag_id - - [VARCHAR(250)] - NOT NULL - -created_at - - [TIMESTAMP] - NOT NULL - -updated_at - - [TIMESTAMP] - NOT NULL + +dag_schedule_dataset_alias_reference + +alias_id + + [INTEGER] + NOT NULL + +dag_id + + [VARCHAR(250)] + NOT NULL + +created_at + + [TIMESTAMP] + NOT NULL + +updated_at + + [TIMESTAMP] + NOT NULL dataset_alias--dag_schedule_dataset_alias_reference - -0..N -1 + +0..N +1 - + dataset - -dataset - -id - - [INTEGER] - NOT NULL - -created_at - - [TIMESTAMP] - NOT NULL - -extra - - [JSON] - NOT NULL - -group - - [VARCHAR(1500)] - NOT NULL - -is_orphaned - - [BOOLEAN] - NOT NULL - -name - - [VARCHAR(1500)] - NOT NULL - -updated_at - - [TIMESTAMP] - NOT NULL - -uri - - [VARCHAR(1500)] - NOT NULL + +dataset + +id + + [INTEGER] + NOT NULL + +created_at + + [TIMESTAMP] + NOT NULL + +extra + + [JSON] + NOT NULL + +group + + [VARCHAR(1500)] + NOT NULL + +is_orphaned + + [BOOLEAN] + NOT NULL + +name + + [VARCHAR(1500)] + NOT NULL + +updated_at + + [TIMESTAMP] + NOT NULL + +uri + + [VARCHAR(1500)] + NOT NULL dataset--dataset_alias_dataset - -0..N -1 + +0..N +1 dataset--dataset_alias_dataset - -0..N -1 + +0..N +1 - + dag_schedule_dataset_reference - -dag_schedule_dataset_reference - -dag_id - - [VARCHAR(250)] - NOT NULL - -dataset_id - - [INTEGER] - NOT NULL - -created_at - - [TIMESTAMP] - NOT NULL - -updated_at - - [TIMESTAMP] - NOT NULL + +dag_schedule_dataset_reference + +dag_id + + [VARCHAR(250)] + NOT NULL + +dataset_id + + [INTEGER] + NOT NULL + +created_at + + [TIMESTAMP] + NOT NULL + +updated_at + + [TIMESTAMP] + NOT NULL dataset--dag_schedule_dataset_reference - -0..N -1 + +0..N +1 - + task_outlet_dataset_reference - -task_outlet_dataset_reference - -dag_id - - [VARCHAR(250)] - NOT NULL - -dataset_id - - [INTEGER] - NOT NULL - -task_id - - [VARCHAR(250)] - NOT NULL - -created_at - - [TIMESTAMP] - NOT NULL - -updated_at - - [TIMESTAMP] - NOT NULL + +task_outlet_dataset_reference + +dag_id + + [VARCHAR(250)] + NOT NULL + +dataset_id + + [INTEGER] + NOT NULL + +task_id + + [VARCHAR(250)] + NOT NULL + +created_at + + [TIMESTAMP] + NOT NULL + +updated_at + + [TIMESTAMP] + NOT NULL dataset--task_outlet_dataset_reference - -0..N -1 + +0..N +1 - + dataset_dag_run_queue - -dataset_dag_run_queue - -dataset_id - - [INTEGER] - NOT NULL - -target_dag_id - - [VARCHAR(250)] - NOT NULL - -created_at - - [TIMESTAMP] - NOT NULL + +dataset_dag_run_queue + +dataset_id + + [INTEGER] + NOT NULL + +target_dag_id + + [VARCHAR(250)] + NOT NULL + +created_at + + [TIMESTAMP] + NOT NULL dataset--dataset_dag_run_queue - -0..N -1 + +0..N +1 - + dataset_event - -dataset_event - -id - - [INTEGER] - NOT NULL - -dataset_id - - [INTEGER] - NOT NULL - -extra - - [JSON] - NOT NULL - -source_dag_id - - [VARCHAR(250)] - -source_map_index - - [INTEGER] - -source_run_id - - [VARCHAR(250)] - -source_task_id - - [VARCHAR(250)] - -timestamp - - [TIMESTAMP] - NOT NULL + +dataset_event + +id + + [INTEGER] + NOT NULL + +dataset_id + + [INTEGER] + NOT NULL + +extra + + [JSON] + NOT NULL + +source_dag_id + + [VARCHAR(250)] + +source_map_index + + [INTEGER] + +source_run_id + + [VARCHAR(250)] + +source_task_id + + [VARCHAR(250)] + +timestamp + + [TIMESTAMP] + NOT NULL dataset_event--dataset_alias_dataset_event - -0..N -1 + +0..N +1 dataset_event--dataset_alias_dataset_event - -0..N -1 + +0..N +1 - + dagrun_dataset_event - -dagrun_dataset_event - -dag_run_id - - [INTEGER] - NOT NULL - -event_id - - [INTEGER] - NOT NULL + +dagrun_dataset_event + +dag_run_id + + [INTEGER] + NOT NULL + +event_id + + [INTEGER] + NOT NULL dataset_event--dagrun_dataset_event - -0..N -1 + +0..N +1 - + dag - -dag - -dag_id - - [VARCHAR(250)] - NOT NULL - -dag_display_name - - [VARCHAR(2000)] - -dataset_expression - - [JSON] - -default_view - - [VARCHAR(25)] - -description - - [TEXT] - -fileloc - - [VARCHAR(2000)] - -has_import_errors - - [BOOLEAN] - -has_task_concurrency_limits - - [BOOLEAN] - NOT NULL - -is_active - - [BOOLEAN] - -is_paused - - [BOOLEAN] - -last_expired - - [TIMESTAMP] - -last_parsed_time - - [TIMESTAMP] - -last_pickled - - [TIMESTAMP] - -max_active_runs - - [INTEGER] - -max_active_tasks - - [INTEGER] - NOT NULL - -max_consecutive_failed_dag_runs - - [INTEGER] - NOT NULL - -next_dagrun - - [TIMESTAMP] - -next_dagrun_create_after - - [TIMESTAMP] - -next_dagrun_data_interval_end - - [TIMESTAMP] - -next_dagrun_data_interval_start - - [TIMESTAMP] - -owners - - [VARCHAR(2000)] - -pickle_id - - [INTEGER] - -processor_subdir - - [VARCHAR(2000)] - -scheduler_lock - - [BOOLEAN] - -timetable_description - - [VARCHAR(1000)] - -timetable_summary - - [TEXT] + +dag + +dag_id + + [VARCHAR(250)] + NOT NULL + +dag_display_name + + [VARCHAR(2000)] + +dataset_expression + + [JSON] + +default_view + + [VARCHAR(25)] + +description + + [TEXT] + +fileloc + + [VARCHAR(2000)] + +has_import_errors + + [BOOLEAN] + +has_task_concurrency_limits + + [BOOLEAN] + NOT NULL + +is_active + + [BOOLEAN] + +is_paused + + [BOOLEAN] + +last_expired + + [TIMESTAMP] + +last_parsed_time + + [TIMESTAMP] + +last_pickled + + [TIMESTAMP] + +max_active_runs + + [INTEGER] + +max_active_tasks + + [INTEGER] + NOT NULL + +max_consecutive_failed_dag_runs + + [INTEGER] + NOT NULL + +next_dagrun + + [TIMESTAMP] + +next_dagrun_create_after + + [TIMESTAMP] + +next_dagrun_data_interval_end + + [TIMESTAMP] + +next_dagrun_data_interval_start + + [TIMESTAMP] + +owners + + [VARCHAR(2000)] + +pickle_id + + [INTEGER] + +processor_subdir + + [VARCHAR(2000)] + +scheduler_lock + + [BOOLEAN] + +timetable_description + + [VARCHAR(1000)] + +timetable_summary + + [TEXT] dag--dag_schedule_dataset_alias_reference - -0..N -1 + +0..N +1 dag--dag_schedule_dataset_reference - -0..N -1 + +0..N +1 dag--task_outlet_dataset_reference - -0..N -1 + +0..N +1 dag--dataset_dag_run_queue - -0..N -1 + +0..N +1 - + dag_tag - -dag_tag - -dag_id - - [VARCHAR(250)] - NOT NULL - -name - - [VARCHAR(100)] - NOT NULL + +dag_tag + +dag_id + + [VARCHAR(250)] + NOT NULL + +name + + [VARCHAR(100)] + NOT NULL dag--dag_tag - -0..N -1 + +0..N +1 - + dag_owner_attributes - -dag_owner_attributes - -dag_id - - [VARCHAR(250)] - NOT NULL - -owner - - [VARCHAR(500)] - NOT NULL - -link - - [VARCHAR(500)] - NOT NULL + +dag_owner_attributes + +dag_id + + [VARCHAR(250)] + NOT NULL + +owner + + [VARCHAR(500)] + NOT NULL + +link + + [VARCHAR(500)] + NOT NULL dag--dag_owner_attributes - -0..N -1 + +0..N +1 - + dag_warning - -dag_warning - -dag_id - - [VARCHAR(250)] - NOT NULL - -warning_type - - [VARCHAR(50)] - NOT NULL - -message - - [TEXT] - NOT NULL - -timestamp - - [TIMESTAMP] - NOT NULL + +dag_warning + +dag_id + + [VARCHAR(250)] + NOT NULL + +warning_type + + [VARCHAR(50)] + NOT NULL + +message + + [TEXT] + NOT NULL + +timestamp + + [TIMESTAMP] + NOT NULL dag--dag_warning - -0..N -1 + +0..N +1 - + log_template - -log_template - -id - - [INTEGER] - NOT NULL - -created_at - - [TIMESTAMP] - NOT NULL - -elasticsearch_id - - [TEXT] - NOT NULL - -filename - - [TEXT] - NOT NULL + +log_template + +id + + [INTEGER] + NOT NULL + +created_at + + [TIMESTAMP] + NOT NULL + +elasticsearch_id + + [TEXT] + NOT NULL + +filename + + [TEXT] + NOT NULL - + dag_run - -dag_run - -id - - [INTEGER] - NOT NULL - -clear_number - - [INTEGER] - NOT NULL - -conf - - [BYTEA] - -creating_job_id - - [INTEGER] - -dag_hash - - [VARCHAR(32)] - -dag_id - - [VARCHAR(250)] - NOT NULL - -data_interval_end - - [TIMESTAMP] - -data_interval_start - - [TIMESTAMP] - -end_date - - [TIMESTAMP] - -external_trigger - - [BOOLEAN] - -last_scheduling_decision - - [TIMESTAMP] - -log_template_id - - [INTEGER] - -logical_date - - [TIMESTAMP] - NOT NULL - -queued_at - - [TIMESTAMP] - -run_id - - [VARCHAR(250)] - NOT NULL + +dag_run -run_type - - [VARCHAR(50)] - NOT NULL +id + + [INTEGER] + NOT NULL -start_date - - [TIMESTAMP] +clear_number + + [INTEGER] + NOT NULL -state - - [VARCHAR(50)] +conf + + [BYTEA] -triggered_by - - [VARCHAR(50)] +creating_job_id + + [INTEGER] -updated_at - - [TIMESTAMP] +dag_id + + [VARCHAR(250)] + NOT NULL + +data_interval_end + + [TIMESTAMP] + +data_interval_start + + [TIMESTAMP] + +end_date + + [TIMESTAMP] + +external_trigger + + [BOOLEAN] + +last_scheduling_decision + + [TIMESTAMP] + +log_template_id + + [INTEGER] + +logical_date + + [TIMESTAMP] + NOT NULL + +queued_at + + [TIMESTAMP] + +run_id + + [VARCHAR(250)] + NOT NULL + +run_type + + [VARCHAR(50)] + NOT NULL + +serialized_dag_id + + [INTEGER] + +start_date + + [TIMESTAMP] + +state + + [VARCHAR(50)] + +triggered_by + + [VARCHAR(50)] + +updated_at + + [TIMESTAMP] log_template--dag_run - -0..N -{0,1} + +0..N +{0,1} dag_run--dagrun_dataset_event - -0..N -1 + +0..N +1 - + task_instance - -task_instance - -dag_id - - [VARCHAR(250)] - NOT NULL - -map_index - - [INTEGER] - NOT NULL - -run_id - - [VARCHAR(250)] - NOT NULL - -task_id - - [VARCHAR(250)] - NOT NULL - -custom_operator_name - - [VARCHAR(1000)] - -duration - - [DOUBLE_PRECISION] - -end_date - - [TIMESTAMP] - -executor - - [VARCHAR(1000)] - -executor_config - - [BYTEA] - -external_executor_id - - [VARCHAR(250)] - -hostname - - [VARCHAR(1000)] - -job_id - - [INTEGER] - -max_tries - - [INTEGER] - -next_kwargs - - [JSON] - -next_method - - [VARCHAR(1000)] - -operator - - [VARCHAR(1000)] - -pid - - [INTEGER] - -pool - - [VARCHAR(256)] - NOT NULL - -pool_slots - - [INTEGER] - NOT NULL - -priority_weight - - [INTEGER] - -queue - - [VARCHAR(256)] - -queued_by_job_id - - [INTEGER] - -queued_dttm - - [TIMESTAMP] - -rendered_map_index - - [VARCHAR(250)] - -start_date - - [TIMESTAMP] - -state - - [VARCHAR(20)] - -task_display_name - - [VARCHAR(2000)] - -trigger_id - - [INTEGER] - -trigger_timeout - - [TIMESTAMP] - -try_number - - [INTEGER] - -unixname - - [VARCHAR(1000)] - -updated_at - - [TIMESTAMP] + +task_instance + +dag_id + + [VARCHAR(250)] + NOT NULL + +map_index + + [INTEGER] + NOT NULL + +run_id + + [VARCHAR(250)] + NOT NULL + +task_id + + [VARCHAR(250)] + NOT NULL + +custom_operator_name + + [VARCHAR(1000)] + +duration + + [DOUBLE_PRECISION] + +end_date + + [TIMESTAMP] + +executor + + [VARCHAR(1000)] + +executor_config + + [BYTEA] + +external_executor_id + + [VARCHAR(250)] + +hostname + + [VARCHAR(1000)] + +job_id + + [INTEGER] + +max_tries + + [INTEGER] + +next_kwargs + + [JSON] + +next_method + + [VARCHAR(1000)] + +operator + + [VARCHAR(1000)] + +pid + + [INTEGER] + +pool + + [VARCHAR(256)] + NOT NULL + +pool_slots + + [INTEGER] + NOT NULL + +priority_weight + + [INTEGER] + +queue + + [VARCHAR(256)] + +queued_by_job_id + + [INTEGER] + +queued_dttm + + [TIMESTAMP] + +rendered_map_index + + [VARCHAR(250)] + +serialized_dag_id + + [INTEGER] + +start_date + + [TIMESTAMP] + +state + + [VARCHAR(20)] + +task_display_name + + [VARCHAR(2000)] + +trigger_id + + [INTEGER] + +trigger_timeout + + [TIMESTAMP] + +try_number + + [INTEGER] + +unixname + + [VARCHAR(1000)] + +updated_at + + [TIMESTAMP] dag_run--task_instance - -0..N -1 + +0..N +1 dag_run--task_instance - -0..N -1 + +0..N +1 - + dag_run_note - -dag_run_note - -dag_run_id - - [INTEGER] - NOT NULL - -content - - [VARCHAR(1000)] - -created_at - - [TIMESTAMP] - NOT NULL - -updated_at - - [TIMESTAMP] - NOT NULL - -user_id - - [VARCHAR(100)] + +dag_run_note + +dag_run_id + + [INTEGER] + NOT NULL + +content + + [VARCHAR(1000)] + +created_at + + [TIMESTAMP] + NOT NULL + +updated_at + + [TIMESTAMP] + NOT NULL + +user_id + + [VARCHAR(128)] dag_run--dag_run_note - -1 -1 + +1 +1 + + + +backfill_dag_run + +backfill_dag_run + +id + + [INTEGER] + NOT NULL + +backfill_id + + [INTEGER] + NOT NULL + +dag_run_id + + [INTEGER] + +sort_ordinal + + [INTEGER] + NOT NULL + + + +dag_run--backfill_dag_run + +0..N +{0,1} - + task_reschedule - -task_reschedule - -id - - [INTEGER] - NOT NULL - -dag_id - - [VARCHAR(250)] - NOT NULL - -duration - - [INTEGER] - NOT NULL - -end_date - - [TIMESTAMP] - NOT NULL - -map_index - - [INTEGER] - NOT NULL - -reschedule_date - - [TIMESTAMP] - NOT NULL - -run_id - - [VARCHAR(250)] - NOT NULL - -start_date - - [TIMESTAMP] - NOT NULL - -task_id - - [VARCHAR(250)] - NOT NULL - -try_number - - [INTEGER] - NOT NULL + +task_reschedule + +id + + [INTEGER] + NOT NULL + +dag_id + + [VARCHAR(250)] + NOT NULL + +duration + + [INTEGER] + NOT NULL + +end_date + + [TIMESTAMP] + NOT NULL + +map_index + + [INTEGER] + NOT NULL + +reschedule_date + + [TIMESTAMP] + NOT NULL + +run_id + + [VARCHAR(250)] + NOT NULL + +start_date + + [TIMESTAMP] + NOT NULL + +task_id + + [VARCHAR(250)] + NOT NULL + +try_number + + [INTEGER] + NOT NULL - + dag_run--task_reschedule - -0..N -1 + +0..N +1 - + dag_run--task_reschedule - -0..N -1 + +0..N +1 - + task_instance--task_reschedule - -0..N -1 + +0..N +1 - + task_instance--task_reschedule - -0..N -1 + +0..N +1 - + task_instance--task_reschedule - -0..N -1 + +0..N +1 - + task_instance--task_reschedule - -0..N -1 + +0..N +1 - + rendered_task_instance_fields - -rendered_task_instance_fields - -dag_id - - [VARCHAR(250)] - NOT NULL - -map_index - - [INTEGER] - NOT NULL - -run_id - - [VARCHAR(250)] - NOT NULL - -task_id - - [VARCHAR(250)] - NOT NULL - -k8s_pod_yaml - - [JSON] - -rendered_fields - - [JSON] - NOT NULL + +rendered_task_instance_fields + +dag_id + + [VARCHAR(250)] + NOT NULL + +map_index + + [INTEGER] + NOT NULL + +run_id + + [VARCHAR(250)] + NOT NULL + +task_id + + [VARCHAR(250)] + NOT NULL + +k8s_pod_yaml + + [JSON] + +rendered_fields + + [JSON] + NOT NULL - + task_instance--rendered_task_instance_fields - -0..N -1 + +0..N +1 - + task_instance--rendered_task_instance_fields - -0..N -1 + +0..N +1 - + task_instance--rendered_task_instance_fields - -0..N -1 + +0..N +1 - + task_instance--rendered_task_instance_fields - -0..N -1 + +0..N +1 - + task_fail - -task_fail - -id - - [INTEGER] - NOT NULL - -dag_id - - [VARCHAR(250)] - NOT NULL - -duration - - [INTEGER] - -end_date - - [TIMESTAMP] - -map_index - - [INTEGER] - NOT NULL - -run_id - - [VARCHAR(250)] - NOT NULL - -start_date - - [TIMESTAMP] - -task_id - - [VARCHAR(250)] - NOT NULL + +task_fail + +id + + [INTEGER] + NOT NULL + +dag_id + + [VARCHAR(250)] + NOT NULL + +duration + + [INTEGER] + +end_date + + [TIMESTAMP] + +map_index + + [INTEGER] + NOT NULL + +run_id + + [VARCHAR(250)] + NOT NULL + +start_date + + [TIMESTAMP] + +task_id + + [VARCHAR(250)] + NOT NULL - + task_instance--task_fail - -0..N -1 + +0..N +1 - + task_instance--task_fail - -0..N -1 + +0..N +1 - + task_instance--task_fail - -0..N -1 + +0..N +1 - + task_instance--task_fail - -0..N -1 + +0..N +1 - + task_map - -task_map - -dag_id - - [VARCHAR(250)] - NOT NULL - -map_index - - [INTEGER] - NOT NULL - -run_id - - [VARCHAR(250)] - NOT NULL - -task_id - - [VARCHAR(250)] - NOT NULL - -keys - - [JSON] - -length - - [INTEGER] - NOT NULL + +task_map + +dag_id + + [VARCHAR(250)] + NOT NULL + +map_index + + [INTEGER] + NOT NULL + +run_id + + [VARCHAR(250)] + NOT NULL + +task_id + + [VARCHAR(250)] + NOT NULL + +keys + + [JSON] + +length + + [INTEGER] + NOT NULL - + task_instance--task_map - -0..N -1 + +0..N +1 - + task_instance--task_map - -0..N -1 + +0..N +1 - + task_instance--task_map - -0..N -1 + +0..N +1 - + task_instance--task_map - -0..N -1 + +0..N +1 - + xcom - -xcom - -dag_run_id - - [INTEGER] - NOT NULL - -key - - [VARCHAR(512)] - NOT NULL - -map_index - - [INTEGER] - NOT NULL - -task_id - - [VARCHAR(250)] - NOT NULL - -dag_id - - [VARCHAR(250)] - NOT NULL - -run_id - - [VARCHAR(250)] - NOT NULL - -timestamp - - [TIMESTAMP] - NOT NULL - -value - - [BYTEA] + +xcom + +dag_run_id + + [INTEGER] + NOT NULL + +key + + [VARCHAR(512)] + NOT NULL + +map_index + + [INTEGER] + NOT NULL + +task_id + + [VARCHAR(250)] + NOT NULL + +dag_id + + [VARCHAR(250)] + NOT NULL + +run_id + + [VARCHAR(250)] + NOT NULL + +timestamp + + [TIMESTAMP] + NOT NULL + +value + + [BYTEA] - + task_instance--xcom - -0..N -1 + +0..N +1 - + task_instance--xcom - -0..N -1 + +0..N +1 - + task_instance--xcom - -0..N -1 + +0..N +1 - + task_instance--xcom - -0..N -1 + +0..N +1 - + task_instance_note - -task_instance_note - -dag_id - - [VARCHAR(250)] - NOT NULL - -map_index - - [INTEGER] - NOT NULL - -run_id - - [VARCHAR(250)] - NOT NULL - -task_id - - [VARCHAR(250)] - NOT NULL - -content - - [VARCHAR(1000)] - -created_at - - [TIMESTAMP] - NOT NULL - -updated_at - - [TIMESTAMP] - NOT NULL - -user_id - - [VARCHAR(100)] + +task_instance_note + +dag_id + + [VARCHAR(250)] + NOT NULL + +map_index + + [INTEGER] + NOT NULL + +run_id + + [VARCHAR(250)] + NOT NULL + +task_id + + [VARCHAR(250)] + NOT NULL + +content + + [VARCHAR(1000)] + +created_at + + [TIMESTAMP] + NOT NULL + +updated_at + + [TIMESTAMP] + NOT NULL + +user_id + + [VARCHAR(128)] - + task_instance--task_instance_note - -0..N -1 + +0..N +1 - + task_instance--task_instance_note - -0..N -1 + +0..N +1 - + task_instance--task_instance_note - -0..N -1 + +0..N +1 - + task_instance--task_instance_note - -0..N -1 + +0..N +1 - + task_instance_history - -task_instance_history - -id - - [INTEGER] - NOT NULL - -custom_operator_name - - [VARCHAR(1000)] - -dag_id - - [VARCHAR(250)] - NOT NULL - -duration - - [DOUBLE_PRECISION] - -end_date - - [TIMESTAMP] - -executor - - [VARCHAR(1000)] - -executor_config - - [BYTEA] - -external_executor_id - - [VARCHAR(250)] - -hostname - - [VARCHAR(1000)] - -job_id - - [INTEGER] - -map_index - - [INTEGER] - NOT NULL - -max_tries - - [INTEGER] - -next_kwargs - - [JSON] - -next_method - - [VARCHAR(1000)] - -operator - - [VARCHAR(1000)] - -pid - - [INTEGER] - -pool - - [VARCHAR(256)] - NOT NULL - -pool_slots - - [INTEGER] - NOT NULL - -priority_weight - - [INTEGER] - -queue - - [VARCHAR(256)] - -queued_by_job_id - - [INTEGER] - -queued_dttm - - [TIMESTAMP] - -rendered_map_index - - [VARCHAR(250)] - -run_id - - [VARCHAR(250)] - NOT NULL - -start_date - - [TIMESTAMP] - -state - - [VARCHAR(20)] - -task_display_name - - [VARCHAR(2000)] - -task_id - - [VARCHAR(250)] - NOT NULL - -trigger_id - - [INTEGER] - -trigger_timeout - - [TIMESTAMP] - -try_number - - [INTEGER] - NOT NULL - -unixname - - [VARCHAR(1000)] - -updated_at - - [TIMESTAMP] + +task_instance_history + +id + + [INTEGER] + NOT NULL + +custom_operator_name + + [VARCHAR(1000)] + +dag_id + + [VARCHAR(250)] + NOT NULL + +duration + + [DOUBLE_PRECISION] + +end_date + + [TIMESTAMP] + +executor + + [VARCHAR(1000)] + +executor_config + + [BYTEA] + +external_executor_id + + [VARCHAR(250)] + +hostname + + [VARCHAR(1000)] + +job_id + + [INTEGER] + +map_index + + [INTEGER] + NOT NULL + +max_tries + + [INTEGER] + +next_kwargs + + [JSON] + +next_method + + [VARCHAR(1000)] + +operator + + [VARCHAR(1000)] + +pid + + [INTEGER] + +pool + + [VARCHAR(256)] + NOT NULL + +pool_slots + + [INTEGER] + NOT NULL + +priority_weight + + [INTEGER] + +queue + + [VARCHAR(256)] + +queued_by_job_id + + [INTEGER] + +queued_dttm + + [TIMESTAMP] + +rendered_map_index + + [VARCHAR(250)] + +run_id + + [VARCHAR(250)] + NOT NULL + +serialized_dag_id + + [INTEGER] + +start_date + + [TIMESTAMP] + +state + + [VARCHAR(20)] + +task_display_name + + [VARCHAR(2000)] + +task_id + + [VARCHAR(250)] + NOT NULL + +trigger_id + + [INTEGER] + +trigger_timeout + + [TIMESTAMP] + +try_number + + [INTEGER] + NOT NULL + +unixname + + [VARCHAR(1000)] + +updated_at + + [TIMESTAMP] - + task_instance--task_instance_history - -0..N -1 + +0..N +1 - + task_instance--task_instance_history - -0..N -1 + +0..N +1 - + task_instance--task_instance_history - -0..N -1 + +0..N +1 - + task_instance--task_instance_history - -0..N -1 + +0..N +1 + + + +serialized_dag + +serialized_dag + +id + + [INTEGER] + NOT NULL + +dag_hash + + [VARCHAR(32)] + NOT NULL + +dag_id + + [VARCHAR(250)] + +data + + [JSON] + +data_compressed + + [BYTEA] + +fileloc + + [VARCHAR(2000)] + NOT NULL + +fileloc_hash + + [BIGINT] + NOT NULL + +last_updated + + [TIMESTAMP] + NOT NULL + +processor_subdir + + [VARCHAR(2000)] + +version_number + + [INTEGER] + NOT NULL + + + +serialized_dag--dag_run + +0..N +{0,1} + + + +serialized_dag--task_instance + +0..N +{0,1} - + trigger - -trigger - -id - - [INTEGER] - NOT NULL - -classpath - - [VARCHAR(1000)] - NOT NULL - -created_date - - [TIMESTAMP] - NOT NULL - -kwargs - - [TEXT] - NOT NULL - -triggerer_id - - [INTEGER] + +trigger + +id + + [INTEGER] + NOT NULL + +classpath + + [VARCHAR(1000)] + NOT NULL + +created_date + + [TIMESTAMP] + NOT NULL + +kwargs + + [TEXT] + NOT NULL + +triggerer_id + + [INTEGER] - + trigger--task_instance - -0..N -{0,1} + +0..N +{0,1} + + + +backfill + +backfill + +id + + [INTEGER] + NOT NULL + +completed_at + + [TIMESTAMP] + +created_at + + [TIMESTAMP] + NOT NULL + +dag_id + + [VARCHAR(250)] + NOT NULL + +dag_run_conf + + [JSON] + +from_date + + [TIMESTAMP] + NOT NULL + +is_paused + + [BOOLEAN] + +max_active_runs + + [INTEGER] + NOT NULL + +to_date + + [TIMESTAMP] + NOT NULL + +updated_at + + [TIMESTAMP] + NOT NULL + + + +backfill--backfill_dag_run + +0..N +1 session - -session - -id - - [INTEGER] - NOT NULL - -data - - [BYTEA] - -expiry - - [TIMESTAMP] - -session_id - - [VARCHAR(255)] + +session + +id + + [INTEGER] + NOT NULL + +data + + [BYTEA] + +expiry + + [TIMESTAMP] + +session_id + + [VARCHAR(255)] alembic_version - -alembic_version - -version_num - - [VARCHAR(32)] - NOT NULL + +alembic_version + +version_num + + [VARCHAR(32)] + NOT NULL ab_user - -ab_user - -id - - [INTEGER] - NOT NULL - -active - - [BOOLEAN] - -changed_by_fk - - [INTEGER] - -changed_on - - [TIMESTAMP] - -created_by_fk - - [INTEGER] - -created_on - - [TIMESTAMP] - -email - - [VARCHAR(512)] - NOT NULL - -fail_login_count - - [INTEGER] - -first_name - - [VARCHAR(256)] - NOT NULL - -last_login - - [TIMESTAMP] - -last_name - - [VARCHAR(256)] - NOT NULL - -login_count - - [INTEGER] - -password - - [VARCHAR(256)] - -username - - [VARCHAR(512)] - NOT NULL + +ab_user + +id + + [INTEGER] + NOT NULL + +active + + [BOOLEAN] + +changed_by_fk + + [INTEGER] + +changed_on + + [TIMESTAMP] + +created_by_fk + + [INTEGER] + +created_on + + [TIMESTAMP] + +email + + [VARCHAR(512)] + NOT NULL + +fail_login_count + + [INTEGER] + +first_name + + [VARCHAR(256)] + NOT NULL + +last_login + + [TIMESTAMP] + +last_name + + [VARCHAR(256)] + NOT NULL + +login_count + + [INTEGER] + +password + + [VARCHAR(256)] + +username + + [VARCHAR(512)] + NOT NULL - + ab_user--ab_user - -0..N -{0,1} + +0..N +{0,1} - + ab_user--ab_user - -0..N -{0,1} + +0..N +{0,1} ab_user_role - -ab_user_role - -id - - [INTEGER] - NOT NULL - -role_id - - [INTEGER] - -user_id - - [INTEGER] + +ab_user_role + +id + + [INTEGER] + NOT NULL + +role_id + + [INTEGER] + +user_id + + [INTEGER] - + ab_user--ab_user_role - -0..N -{0,1} + +0..N +{0,1} ab_register_user - -ab_register_user - -id - - [INTEGER] - NOT NULL - -email - - [VARCHAR(512)] - NOT NULL - -first_name - - [VARCHAR(256)] - NOT NULL - -last_name - - [VARCHAR(256)] - NOT NULL - -password - - [VARCHAR(256)] - -registration_date - - [TIMESTAMP] - -registration_hash - - [VARCHAR(256)] - -username - - [VARCHAR(512)] - NOT NULL + +ab_register_user + +id + + [INTEGER] + NOT NULL + +email + + [VARCHAR(512)] + NOT NULL + +first_name + + [VARCHAR(256)] + NOT NULL + +last_name + + [VARCHAR(256)] + NOT NULL + +password + + [VARCHAR(256)] + +registration_date + + [TIMESTAMP] + +registration_hash + + [VARCHAR(256)] + +username + + [VARCHAR(512)] + NOT NULL ab_permission - -ab_permission - -id - - [INTEGER] - NOT NULL - -name - - [VARCHAR(100)] - NOT NULL + +ab_permission + +id + + [INTEGER] + NOT NULL + +name + + [VARCHAR(100)] + NOT NULL ab_permission_view - -ab_permission_view - -id - - [INTEGER] - NOT NULL - -permission_id - - [INTEGER] - -view_menu_id - - [INTEGER] + +ab_permission_view + +id + + [INTEGER] + NOT NULL + +permission_id + + [INTEGER] + +view_menu_id + + [INTEGER] - + ab_permission--ab_permission_view - -0..N -{0,1} + +0..N +{0,1} ab_permission_view_role - -ab_permission_view_role - -id - - [INTEGER] - NOT NULL - -permission_view_id - - [INTEGER] - -role_id - - [INTEGER] + +ab_permission_view_role + +id + + [INTEGER] + NOT NULL + +permission_view_id + + [INTEGER] + +role_id + + [INTEGER] - + ab_permission_view--ab_permission_view_role - -0..N -{0,1} + +0..N +{0,1} ab_view_menu - -ab_view_menu - -id - - [INTEGER] - NOT NULL - -name - - [VARCHAR(250)] - NOT NULL + +ab_view_menu + +id + + [INTEGER] + NOT NULL + +name + + [VARCHAR(250)] + NOT NULL - + ab_view_menu--ab_permission_view - -0..N -{0,1} + +0..N +{0,1} ab_role - -ab_role - -id - - [INTEGER] - NOT NULL - -name - - [VARCHAR(64)] - NOT NULL + +ab_role + +id + + [INTEGER] + NOT NULL + +name + + [VARCHAR(64)] + NOT NULL - + ab_role--ab_user_role - -0..N -{0,1} + +0..N +{0,1} - + ab_role--ab_permission_view_role - -0..N -{0,1} + +0..N +{0,1} alembic_version_fab - -alembic_version_fab - -version_num - - [VARCHAR(32)] - NOT NULL + +alembic_version_fab + +version_num + + [VARCHAR(32)] + NOT NULL diff --git a/docs/apache-airflow/migrations-ref.rst b/docs/apache-airflow/migrations-ref.rst index e4fb2dfa332eb..067cbb6f619e6 100644 --- a/docs/apache-airflow/migrations-ref.rst +++ b/docs/apache-airflow/migrations-ref.rst @@ -39,7 +39,11 @@ Here's the list of all the Database Migrations that are executed via when you ru +-------------------------+------------------+-------------------+--------------------------------------------------------------+ | Revision ID | Revises ID | Airflow Version | Description | +=========================+==================+===================+==============================================================+ -| ``0d9e73a75ee4`` (head) | ``44eabb1904b4`` | ``3.0.0`` | Add name and group fields to DatasetModel. | +| ``4235395d5ec5`` (head) | ``e1ff90d3efe9`` | ``3.0.0`` | Add SDM foreign key to DagRun, TI & TIH. | ++-------------------------+------------------+-------------------+--------------------------------------------------------------+ +| ``e1ff90d3efe9`` | ``0d9e73a75ee4`` | ``3.0.0`` | Add serial ID to SDM. | ++-------------------------+------------------+-------------------+--------------------------------------------------------------+ +| ``0d9e73a75ee4`` | ``44eabb1904b4`` | ``3.0.0`` | Add name and group fields to DatasetModel. | +-------------------------+------------------+-------------------+--------------------------------------------------------------+ | ``44eabb1904b4`` | ``16cbcb1c8c36`` | ``3.0.0`` | Update dag_run_note.user_id and task_instance_note.user_id | | | | | columns to String. | diff --git a/scripts/ci/pre_commit/check_ti_vs_tis_attributes.py b/scripts/ci/pre_commit/check_ti_vs_tis_attributes.py index 1dfc51a0a040e..457d1b612d11e 100755 --- a/scripts/ci/pre_commit/check_ti_vs_tis_attributes.py +++ b/scripts/ci/pre_commit/check_ti_vs_tis_attributes.py @@ -52,6 +52,7 @@ def compare_attributes(path1, path2): "triggerer_job", "note", "rendered_task_instance_fields", + "serialized_dag", } # exclude attrs not necessary to be in TaskInstanceHistory if not diff: return diff --git a/tests/conftest.py b/tests/conftest.py index 60d009416fe8e..1a490a77c068d 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -986,15 +986,16 @@ def cleanup(self): # To isolate problems here with problems from elsewhere on the session object self.session.rollback() - self.session.query(SerializedDagModel).filter( - SerializedDagModel.dag_id.in_(dag_ids) - ).delete(synchronize_session=False) self.session.query(DagRun).filter(DagRun.dag_id.in_(dag_ids)).delete( synchronize_session=False, ) + self.session.query(TaskInstance).filter(TaskInstance.dag_id.in_(dag_ids)).delete( synchronize_session=False, ) + self.session.query(SerializedDagModel).filter( + SerializedDagModel.dag_id.in_(dag_ids) + ).delete(synchronize_session=False) self.session.query(XCom).filter(XCom.dag_id.in_(dag_ids)).delete( synchronize_session=False, ) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index c8b40f6af7e24..ac7910afbc579 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3471,7 +3471,7 @@ def test_verify_integrity_if_dag_not_changed(self, dag_maker): assert tis_count == 1 latest_dag_version = SerializedDagModel.get_latest_version_hash(dr.dag_id, session=session) - assert dr.dag_hash == latest_dag_version + assert dr.dag_hash() == latest_dag_version session.rollback() session.close() @@ -3505,7 +3505,7 @@ def test_verify_integrity_if_dag_changed(self, dag_maker): dr = drs[0] dag_version_1 = SerializedDagModel.get_latest_version_hash(dr.dag_id, session=session) - assert dr.dag_hash == dag_version_1 + assert dr.dag_hash() == dag_version_1 assert self.job_runner.dagbag.dags == {"test_verify_integrity_if_dag_changed": dag} assert len(self.job_runner.dagbag.dags.get("test_verify_integrity_if_dag_changed").tasks) == 1 @@ -3522,7 +3522,7 @@ def test_verify_integrity_if_dag_changed(self, dag_maker): drs = DagRun.find(dag_id=dag.dag_id, session=session) assert len(drs) == 1 dr = drs[0] - assert dr.dag_hash == dag_version_2 + assert dr.dag_hash() == dag_version_2 assert self.job_runner.dagbag.dags == {"test_verify_integrity_if_dag_changed": dag} assert len(self.job_runner.dagbag.dags.get("test_verify_integrity_if_dag_changed").tasks) == 2 @@ -3538,7 +3538,7 @@ def test_verify_integrity_if_dag_changed(self, dag_maker): assert tis_count == 2 latest_dag_version = SerializedDagModel.get_latest_version_hash(dr.dag_id, session=session) - assert dr.dag_hash == latest_dag_version + assert dr.dag_hash() == latest_dag_version session.rollback() session.close() @@ -3572,7 +3572,7 @@ def test_verify_integrity_if_dag_disappeared(self, dag_maker, caplog): dr = drs[0] dag_version_1 = SerializedDagModel.get_latest_version_hash(dag_id, session=session) - assert dr.dag_hash == dag_version_1 + assert dr.dag_hash() == dag_version_1 assert self.job_runner.dagbag.dags == {"test_verify_integrity_if_dag_disappeared": dag} assert len(self.job_runner.dagbag.dags.get("test_verify_integrity_if_dag_disappeared").tasks) == 1 diff --git a/tests/models/test_serialized_dag.py b/tests/models/test_serialized_dag.py index d9a77e55edaf5..9b9669a3b028a 100644 --- a/tests/models/test_serialized_dag.py +++ b/tests/models/test_serialized_dag.py @@ -69,10 +69,13 @@ def test_dag_fileloc_hash(self): """Verifies the correctness of hashing file path.""" assert DagCode.dag_fileloc_hash("/airflow/dags/test_dag.py") == 33826252060516589 - def _write_example_dags(self): + def _write_example_dags(self, processor_subdir=None): example_dags = make_example_dags(example_dags_module) for dag in example_dags.values(): - SDM.write_dag(dag) + if not processor_subdir: + SDM.write_dag(dag) + else: + SDM.write_dag(dag=dag, processor_subdir=processor_subdir) return example_dags @pytest.mark.skip_if_database_isolation_mode # Does not work in db isolation mode @@ -98,12 +101,12 @@ def test_serialized_dag_is_updated_if_dag_is_changed(self): assert dag_updated is True with create_session() as session: - s_dag = session.get(SDM, example_bash_op_dag.dag_id) + s_dag = session.scalar(SDM.latest_item_select_object(example_bash_op_dag.dag_id)) # Test that if DAG is not changed, Serialized DAG is not re-written and last_updated # column is not updated dag_updated = SDM.write_dag(dag=example_bash_op_dag) - s_dag_1 = session.get(SDM, example_bash_op_dag.dag_id) + s_dag_1 = session.scalar(SDM.latest_item_select_object(example_bash_op_dag.dag_id)) assert s_dag_1.dag_hash == s_dag.dag_hash assert s_dag.last_updated == s_dag_1.last_updated @@ -114,7 +117,7 @@ def test_serialized_dag_is_updated_if_dag_is_changed(self): assert example_bash_op_dag.tags == {"example", "example2", "new_tag"} dag_updated = SDM.write_dag(dag=example_bash_op_dag) - s_dag_2 = session.get(SDM, example_bash_op_dag.dag_id) + s_dag_2 = session.scalar(SDM.latest_item_select_object(example_bash_op_dag.dag_id)) assert s_dag.last_updated != s_dag_2.last_updated assert s_dag.dag_hash != s_dag_2.dag_hash @@ -130,12 +133,12 @@ def test_serialized_dag_is_updated_if_processor_subdir_changed(self): assert dag_updated is True with create_session() as session: - s_dag = session.get(SDM, example_bash_op_dag.dag_id) + s_dag = session.scalar(SDM.latest_item_select_object(example_bash_op_dag.dag_id)) # Test that if DAG is not changed, Serialized DAG is not re-written and last_updated # column is not updated dag_updated = SDM.write_dag(dag=example_bash_op_dag, processor_subdir="/tmp/test") - s_dag_1 = session.get(SDM, example_bash_op_dag.dag_id) + s_dag_1 = session.scalar(SDM.latest_item_select_object(example_bash_op_dag.dag_id)) assert s_dag_1.dag_hash == s_dag.dag_hash assert s_dag.last_updated == s_dag_1.last_updated @@ -144,7 +147,7 @@ def test_serialized_dag_is_updated_if_processor_subdir_changed(self): # Update DAG dag_updated = SDM.write_dag(dag=example_bash_op_dag, processor_subdir="/tmp/other") - s_dag_2 = session.get(SDM, example_bash_op_dag.dag_id) + s_dag_2 = session.scalar(SDM.latest_item_select_object(example_bash_op_dag.dag_id)) assert s_dag.processor_subdir != s_dag_2.processor_subdir assert dag_updated is True @@ -189,7 +192,7 @@ def test_bulk_sync_to_db(self): DAG("dag_2", schedule=None), DAG("dag_3", schedule=None), ] - with assert_queries_count(10): + with assert_queries_count(12): SDM.bulk_sync_to_db(dags) @pytest.mark.skip_if_database_isolation_mode # Does not work in db isolation mode @@ -283,3 +286,20 @@ def get_hash_set(): first_hashes = get_hash_set() # assert that the hashes are the same assert first_hashes == get_hash_set() + + def test_versions_increment_on_serialized_dag_change(self, session): + example_dags = make_example_dags(example_dags_module) + example_bash_op_dag = example_dags.get("example_bash_operator") + for i in range(3): + SDM.write_dag(dag=example_bash_op_dag, processor_subdir=f"/tmp/test-{i+1}") + + sdm_dags = session.scalars(select(SDM).where(SDM.dag_id == "example_bash_operator")).all() + assert len(sdm_dags) == 3 + assert sorted([s.version_number for s in sdm_dags]) == [1, 2, 3] + + def test_read_all_dags_gets_the_latest_of_the_serdags(self, session): + dags = self._write_example_dags() + # sync it again + self._write_example_dags(processor_subdir="/tmp/subdir") + all_dags = SDM.read_all_dags() + assert len(dags) == len(all_dags) diff --git a/tests/models/test_taskinstance.py b/tests/models/test_taskinstance.py index 8c334366f0488..79da0ea357ba5 100644 --- a/tests/models/test_taskinstance.py +++ b/tests/models/test_taskinstance.py @@ -4020,6 +4020,7 @@ def test_refresh_from_db(self, create_task_instance): "next_kwargs": None, "next_method": None, "updated_at": None, + "serialized_dag_id": None, "task_display_name": "Test Refresh from DB Task", } # Make sure we aren't missing any new value in our expected_values list.