diff --git a/airflow/datasets/manager.py b/airflow/datasets/manager.py index 2c3eda4dcc106..7237b7544355b 100644 --- a/airflow/datasets/manager.py +++ b/airflow/datasets/manager.py @@ -28,6 +28,8 @@ from airflow.utils.log.logging_mixin import LoggingMixin if TYPE_CHECKING: + from datetime import datetime + from airflow.models.taskinstance import TaskInstance @@ -53,22 +55,30 @@ def register_dataset_change( if not dataset_model: self.log.warning("DatasetModel %s not found", dataset_model) return - session.add( - DatasetEvent( - dataset_id=dataset_model.id, - source_task_id=task_instance.task_id, - source_dag_id=task_instance.dag_id, - source_run_id=task_instance.run_id, - source_map_index=task_instance.map_index, - extra=extra, - ) + dataset_event = DatasetEvent( + dataset_id=dataset_model.id, + source_task_id=task_instance.task_id, + source_dag_id=task_instance.dag_id, + source_run_id=task_instance.run_id, + source_map_index=task_instance.map_index, + extra=extra, ) + session.add(dataset_event) session.flush() - if dataset_model.consuming_dags: - self._queue_dagruns(dataset_model, session) + + downstream_dag_ids = [x.dag_id for x in dataset_model.consuming_dags] + if downstream_dag_ids: + self._queue_dagruns( + dag_ids=downstream_dag_ids, + dataset_id=dataset_model.id, + event_timestamp=dataset_event.timestamp, + session=session, + ) session.flush() - def _queue_dagruns(self, dataset: DatasetModel, session: Session) -> None: + def _queue_dagruns( + self, *, dag_ids: list[str], dataset_id: int, event_timestamp: datetime, session: Session + ) -> None: # Possible race condition: if multiple dags or multiple (usually # mapped) tasks update the same dataset, this can fail with a unique # constraint violation. @@ -79,31 +89,45 @@ def _queue_dagruns(self, dataset: DatasetModel, session: Session) -> None: # where `ti.state` is changed. if session.bind.dialect.name == "postgresql": - return self._postgres_queue_dagruns(dataset, session) - return self._slow_path_queue_dagruns(dataset, session) + return self._postgres_queue_dagruns( + dag_ids=dag_ids, + dataset_id=dataset_id, + event_timestamp=event_timestamp, + session=session, + ) + return self._slow_path_queue_dagruns( + dag_ids=dag_ids, + dataset_id=dataset_id, + event_timestamp=event_timestamp, + session=session, + ) - def _slow_path_queue_dagruns(self, dataset: DatasetModel, session: Session) -> None: - consuming_dag_ids = [x.dag_id for x in dataset.consuming_dags] - self.log.debug("consuming dag ids %s", consuming_dag_ids) + def _slow_path_queue_dagruns(self, *, dag_ids, dataset_id, event_timestamp, session: Session) -> None: + self.log.debug("consuming dag ids %s", dag_ids) # Don't error whole transaction when a single RunQueue item conflicts. # https://docs.sqlalchemy.org/en/14/orm/session_transaction.html#using-savepoint - for dag_id in consuming_dag_ids: - item = DatasetDagRunQueue(target_dag_id=dag_id, dataset_id=dataset.id) + for dag_id in dag_ids: + item = DatasetDagRunQueue( + target_dag_id=dag_id, + dataset_id=dataset_id, + event_timestamp=event_timestamp, + ) try: with session.begin_nested(): session.merge(item) except exc.IntegrityError: self.log.debug("Skipping record %s", item, exc_info=True) - def _postgres_queue_dagruns(self, dataset: DatasetModel, session: Session) -> None: + def _postgres_queue_dagruns(self, *, dag_ids, dataset_id, event_timestamp, session: Session) -> None: from sqlalchemy.dialects.postgresql import insert - stmt = insert(DatasetDagRunQueue).values(dataset_id=dataset.id).on_conflict_do_nothing() - session.execute( - stmt, - [{'target_dag_id': target_dag.dag_id} for target_dag in dataset.consuming_dags], + stmt = ( + insert(DatasetDagRunQueue) + .values(dataset_id=dataset_id, event_timestamp=event_timestamp) + .on_conflict_do_nothing() ) + session.execute(stmt, [{'target_dag_id': x} for x in dag_ids]) def resolve_dataset_manager() -> DatasetManager: diff --git a/airflow/jobs/scheduler_job.py b/airflow/jobs/scheduler_job.py index 53a96cf9acdbd..ff2050c81c52b 100644 --- a/airflow/jobs/scheduler_job.py +++ b/airflow/jobs/scheduler_job.py @@ -1087,15 +1087,14 @@ def _create_dag_runs_dataset_triggered( """For DAGs that are triggered by datasets, create dag runs.""" # Bulk Fetch DagRuns with dag_id and execution_date same # as DagModel.dag_id and DagModel.next_dagrun - # This list is used to verify if the DagRun already exist so that we don't attempt to create - # duplicate dag runs - exec_dates = { - dag_id: timezone.coerce_datetime(last_time) - for dag_id, (_, last_time) in dataset_triggered_dag_info.items() + # This list is used to verify if the DagRun already exists + # so that we don't attempt to create duplicate dag runs + dag_event_timestamps = { + dag_id: last_event_time for dag_id, (_, last_event_time) in dataset_triggered_dag_info.items() } - existing_dagruns: set[tuple[str, timezone.DateTime]] = set( + existing_dagruns: set[tuple[str, datetime]] = set( session.query(DagRun.dag_id, DagRun.execution_date).filter( - tuple_in_condition((DagRun.dag_id, DagRun.execution_date), exec_dates.items()) + tuple_in_condition((DagRun.dag_id, DagRun.execution_date), dag_event_timestamps.items()) ) ) @@ -1122,14 +1121,14 @@ def _create_dag_runs_dataset_triggered( # we need to set dag.next_dagrun_info if the Dag Run already exists or if we # create a new one. This is so that in the next Scheduling loop we try to create new runs # instead of falling in a loop of Integrity Error. - exec_date = exec_dates[dag.dag_id] - if (dag.dag_id, exec_date) not in existing_dagruns: + last_event_timestamp = dag_event_timestamps[dag.dag_id] + if (dag.dag_id, last_event_timestamp) not in existing_dagruns: previous_dag_run = ( session.query(DagRun) .filter( DagRun.dag_id == dag.dag_id, - DagRun.execution_date < exec_date, + DagRun.execution_date < last_event_timestamp, DagRun.run_type == DagRunType.DATASET_TRIGGERED, ) .order_by(DagRun.execution_date.desc()) @@ -1137,7 +1136,7 @@ def _create_dag_runs_dataset_triggered( ) dataset_event_filters = [ DagScheduleDatasetReference.dag_id == dag.dag_id, - DatasetEvent.timestamp <= exec_date, + DatasetEvent.timestamp <= last_event_timestamp, ] if previous_dag_run: dataset_event_filters.append(DatasetEvent.timestamp > previous_dag_run.execution_date) @@ -1153,10 +1152,13 @@ def _create_dag_runs_dataset_triggered( .all() ) - data_interval = dag.timetable.data_interval_for_events(exec_date, dataset_events) + data_interval = dag.timetable.data_interval_for_events( + last_event_timestamp, # type: ignore + dataset_events, + ) run_id = dag.timetable.generate_run_id( run_type=DagRunType.DATASET_TRIGGERED, - logical_date=exec_date, + logical_date=last_event_timestamp, # type: ignore data_interval=data_interval, session=session, events=dataset_events, @@ -1165,7 +1167,7 @@ def _create_dag_runs_dataset_triggered( dag_run = dag.create_dagrun( run_id=run_id, run_type=DagRunType.DATASET_TRIGGERED, - execution_date=exec_date, + execution_date=last_event_timestamp, data_interval=data_interval, state=DagRunState.QUEUED, external_trigger=False, diff --git a/airflow/migrations/versions/0118_2_4_2_add_event_timestamp_to_ddrq.py b/airflow/migrations/versions/0118_2_4_2_add_event_timestamp_to_ddrq.py new file mode 100644 index 0000000000000..3e2d43fb18ee2 --- /dev/null +++ b/airflow/migrations/versions/0118_2_4_2_add_event_timestamp_to_ddrq.py @@ -0,0 +1,50 @@ +# +# 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 event timestamp to DDRQ + +Revision ID: 2b72b0fd20ef +Revises: ee8d93fcc81e +Create Date: 2022-09-21 21:28:23.896961 + +""" + +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op + +from airflow.migrations.db_types import TIMESTAMP + +# revision identifiers, used by Alembic. +revision = '2b72b0fd20ef' +down_revision = 'ecb43d2a1842' +branch_labels = None +depends_on = None +airflow_version = '2.4.2' + + +def upgrade(): + """Apply Add event timestamp to DDRQ""" + op.add_column('dataset_dag_run_queue', sa.Column('event_timestamp', TIMESTAMP, nullable=False)) + + +def downgrade(): + """Unapply Add event timestamp to DDRQ""" + with op.batch_alter_table('dataset_dag_run_queue', schema=None) as batch_op: + batch_op.drop_column('event_timestamp') diff --git a/airflow/migrations/versions/0118_2_5_0_add_updated_at_to_dagrun_and_ti.py b/airflow/migrations/versions/0119_2_5_0_add_updated_at_to_dagrun_and_ti.py similarity index 98% rename from airflow/migrations/versions/0118_2_5_0_add_updated_at_to_dagrun_and_ti.py rename to airflow/migrations/versions/0119_2_5_0_add_updated_at_to_dagrun_and_ti.py index fc580529658d6..20711f794ec3c 100644 --- a/airflow/migrations/versions/0118_2_5_0_add_updated_at_to_dagrun_and_ti.py +++ b/airflow/migrations/versions/0119_2_5_0_add_updated_at_to_dagrun_and_ti.py @@ -33,7 +33,7 @@ # revision identifiers, used by Alembic. revision = 'ee8d93fcc81e' -down_revision = 'ecb43d2a1842' +down_revision = '2b72b0fd20ef' branch_labels = None depends_on = None airflow_version = '2.5.0' diff --git a/airflow/models/dataset.py b/airflow/models/dataset.py index b1a58e5442d87..b9ac027a4dc88 100644 --- a/airflow/models/dataset.py +++ b/airflow/models/dataset.py @@ -204,6 +204,7 @@ class DatasetDagRunQueue(Base): dataset_id = Column(Integer, primary_key=True, nullable=False) target_dag_id = Column(StringID(), primary_key=True, nullable=False) + event_timestamp = Column(UtcDateTime, nullable=False) created_at = Column(UtcDateTime, default=timezone.utcnow, nullable=False) __tablename__ = "dataset_dag_run_queue" diff --git a/docs/apache-airflow/img/airflow_erd.sha256 b/docs/apache-airflow/img/airflow_erd.sha256 index 2cde8179cbf12..81414b6d7e912 100644 --- a/docs/apache-airflow/img/airflow_erd.sha256 +++ b/docs/apache-airflow/img/airflow_erd.sha256 @@ -1 +1 @@ -5b101dceaef5d9343cddbfab0db6087b646d031eb8a1c3157f79e6cd811e14f4 \ No newline at end of file +25a09e524c0b06c3762da8f08e4850822f9ed40d5d400a634963678b365f3642 \ 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 f113037c85ea5..e06f1bdf3bdab 100644 --- a/docs/apache-airflow/img/airflow_erd.svg +++ b/docs/apache-airflow/img/airflow_erd.svg @@ -4,11 +4,11 @@ - - + + %3 - + ab_permission @@ -26,48 +26,48 @@ 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} @@ -86,53 +86,53 @@ 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_permission_view_role - -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_role--ab_user_role - -0..N -{0,1} + +0..N +{0,1} @@ -172,76 +172,76 @@ ab_user - -ab_user - -id - [INTEGER] - NOT NULL - -active - [BOOLEAN] - -changed_by_fk - [INTEGER] - -changed_on - [DATETIME] - -created_by_fk - [INTEGER] - -created_on - [DATETIME] - -email - [VARCHAR(256)] - NOT NULL - -fail_login_count - [INTEGER] - -first_name - [VARCHAR(64)] - NOT NULL - -last_login - [DATETIME] - -last_name - [VARCHAR(64)] - NOT NULL - -login_count - [INTEGER] - -password - [VARCHAR(256)] - -username - [VARCHAR(256)] - NOT NULL + +ab_user + +id + [INTEGER] + NOT NULL + +active + [BOOLEAN] + +changed_by_fk + [INTEGER] + +changed_on + [DATETIME] + +created_by_fk + [INTEGER] + +created_on + [DATETIME] + +email + [VARCHAR(256)] + NOT NULL + +fail_login_count + [INTEGER] + +first_name + [VARCHAR(64)] + NOT NULL + +last_login + [DATETIME] + +last_name + [VARCHAR(64)] + NOT NULL + +login_count + [INTEGER] + +password + [VARCHAR(256)] + +username + [VARCHAR(256)] + NOT NULL ab_user--ab_user_role - -0..N -{0,1} + +0..N +{0,1} ab_user--ab_user - -0..N -{0,1} + +0..N +{0,1} ab_user--ab_user - -0..N -{0,1} + +0..N +{0,1} @@ -414,760 +414,764 @@ 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 -{0,1} + +0..N +{0,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 dag--dag_schedule_dataset_reference - -0..N -{0,1} + +0..N +{0,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 -{0,1} + +0..N +{0,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 -{0,1} + +0..N +{0,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 + +event_timestamp + [TIMESTAMP] + NOT NULL dag--dataset_dag_run_queue - -0..N -{0,1} + +0..N +{0,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 dag--task_outlet_dataset_reference - -0..N -{0,1} + +0..N +{0,1} dag_code - -dag_code - -fileloc_hash - [BIGINT] - NOT NULL - -fileloc - [VARCHAR(2000)] - NOT NULL - -last_updated - [TIMESTAMP] - NOT NULL - -source_code - [TEXT] - NOT NULL + +dag_code + +fileloc_hash + [BIGINT] + NOT NULL + +fileloc + [VARCHAR(2000)] + NOT NULL + +last_updated + [TIMESTAMP] + NOT NULL + +source_code + [TEXT] + NOT NULL dag_pickle - -dag_pickle - -id - [INTEGER] - NOT NULL - -created_dttm - [TIMESTAMP] - -pickle - [BLOB] - -pickle_hash - [BIGINT] + +dag_pickle + +id + [INTEGER] + NOT NULL + +created_dttm + [TIMESTAMP] + +pickle + [BLOB] + +pickle_hash + [BIGINT] dag_run - -dag_run - -id - [INTEGER] - NOT NULL - -conf - [BLOB] - -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] - -execution_date - [TIMESTAMP] - NOT NULL - -external_trigger - [BOOLEAN] - -last_scheduling_decision - [TIMESTAMP] - -log_template_id - [INTEGER] - -queued_at - [TIMESTAMP] - -run_id - [VARCHAR(250)] - NOT NULL - -run_type - [VARCHAR(50)] - NOT NULL - -start_date - [TIMESTAMP] - -state - [VARCHAR(50)] - -updated_at - [TIMESTAMP] + +dag_run + +id + [INTEGER] + NOT NULL + +conf + [BLOB] + +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] + +execution_date + [TIMESTAMP] + NOT NULL + +external_trigger + [BOOLEAN] + +last_scheduling_decision + [TIMESTAMP] + +log_template_id + [INTEGER] + +queued_at + [TIMESTAMP] + +run_id + [VARCHAR(250)] + NOT NULL + +run_type + [VARCHAR(50)] + NOT NULL + +start_date + [TIMESTAMP] + +state + [VARCHAR(50)] + +updated_at + [TIMESTAMP] 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 dag_run--dagrun_dataset_event - -0..N -{0,1} + +0..N +{0,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 - -duration - [FLOAT] - -end_date - [TIMESTAMP] - -executor_config - [BLOB] - -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] - -start_date - [TIMESTAMP] - -state - [VARCHAR(20)] - -trigger_id - [INTEGER] - -trigger_timeout - [DATETIME] - -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 + +duration + [FLOAT] + +end_date + [TIMESTAMP] + +executor_config + [BLOB] + +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] + +start_date + [TIMESTAMP] + +state + [VARCHAR(20)] + +trigger_id + [INTEGER] + +trigger_timeout + [DATETIME] + +try_number + [INTEGER] + +unixname + [VARCHAR(1000)] + +updated_at + [TIMESTAMP] dag_run--task_instance - -0..N -{0,1} + +0..N +{0,1} dag_run--task_instance - -0..N -{0,1} + +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 -{0,1} + +0..N +{0,1} dag_run--task_reschedule - -0..N -{0,1} + +0..N +{0,1} task_instance--task_reschedule - -0..N -{0,1} + +0..N +{0,1} task_instance--task_reschedule - -0..N -{0,1} + +0..N +{0,1} task_instance--task_reschedule - -0..N -{0,1} + +0..N +{0,1} task_instance--task_reschedule - -0..N -{0,1} + +0..N +{0,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 -{0,1} + +0..N +{0,1} task_instance--rendered_task_instance_fields - -0..N -{0,1} + +0..N +{0,1} task_instance--rendered_task_instance_fields - -0..N -{0,1} + +0..N +{0,1} task_instance--rendered_task_instance_fields - -0..N -{0,1} + +0..N +{0,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 -{0,1} + +0..N +{0,1} task_instance--task_fail - -0..N -{0,1} + +0..N +{0,1} task_instance--task_fail - -0..N -{0,1} + +0..N +{0,1} task_instance--task_fail - -0..N -{0,1} + +0..N +{0,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 -{0,1} + +0..N +{0,1} task_instance--task_map - -0..N -{0,1} + +0..N +{0,1} task_instance--task_map - -0..N -{0,1} + +0..N +{0,1} task_instance--task_map - -0..N -{0,1} + +0..N +{0,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 - [BLOB] + +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 + [BLOB] task_instance--xcom - -0..N -{0,1} + +0..N +{0,1} task_instance--xcom - -0..N -{0,1} + +0..N +{0,1} task_instance--xcom - -0..N -{0,1} + +0..N +{0,1} task_instance--xcom - -0..N -{0,1} + +0..N +{0,1} log_template - -log_template + +log_template + +id + [INTEGER] + NOT NULL -id - [INTEGER] - NOT NULL +created_at + [TIMESTAMP] + NOT NULL -created_at - [TIMESTAMP] - NOT NULL +elasticsearch_id + [TEXT] + NOT NULL -elasticsearch_id - [TEXT] - NOT NULL - -filename - [TEXT] - NOT NULL +filename + [TEXT] + NOT NULL log_template--dag_run - -0..N -{0,1} + +0..N +{0,1} @@ -1198,311 +1202,311 @@ dataset--dag_schedule_dataset_reference - -0..N -{0,1} + +0..N +{0,1} dataset--dataset_dag_run_queue - -0..N -{0,1} + +0..N +{0,1} dataset--task_outlet_dataset_reference - -0..N + +0..N {0,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--dagrun_dataset_event - -0..N -{0,1} + +0..N +{0,1} import_error - -import_error + +import_error + +id + [INTEGER] + NOT NULL -id - [INTEGER] - NOT NULL +filename + [VARCHAR(1024)] -filename - [VARCHAR(1024)] +stacktrace + [TEXT] -stacktrace - [TEXT] - -timestamp - [TIMESTAMP] +timestamp + [TIMESTAMP] job - -job + +job + +id + [INTEGER] + NOT NULL -id - [INTEGER] - NOT NULL +dag_id + [VARCHAR(250)] -dag_id - [VARCHAR(250)] +end_date + [TIMESTAMP] -end_date - [TIMESTAMP] +executor_class + [VARCHAR(500)] -executor_class - [VARCHAR(500)] +hostname + [VARCHAR(500)] -hostname - [VARCHAR(500)] +job_type + [VARCHAR(30)] -job_type - [VARCHAR(30)] +latest_heartbeat + [TIMESTAMP] -latest_heartbeat - [TIMESTAMP] +start_date + [TIMESTAMP] -start_date - [TIMESTAMP] +state + [VARCHAR(20)] -state - [VARCHAR(20)] - -unixname - [VARCHAR(1000)] +unixname + [VARCHAR(1000)] log - -log + +log + +id + [INTEGER] + NOT NULL -id - [INTEGER] - NOT NULL +dag_id + [VARCHAR(250)] -dag_id - [VARCHAR(250)] +dttm + [TIMESTAMP] -dttm - [TIMESTAMP] +event + [VARCHAR(30)] -event - [VARCHAR(30)] +execution_date + [TIMESTAMP] -execution_date - [TIMESTAMP] +extra + [TEXT] -extra - [TEXT] +map_index + [INTEGER] -map_index - [INTEGER] +owner + [VARCHAR(500)] -owner - [VARCHAR(500)] - -task_id - [VARCHAR(250)] +task_id + [VARCHAR(250)] trigger - -trigger - -id - [INTEGER] - NOT NULL - -classpath - [VARCHAR(1000)] - NOT NULL - -created_date - [TIMESTAMP] - NOT NULL - -kwargs - [JSON] - NOT NULL - -triggerer_id - [INTEGER] + +trigger + +id + [INTEGER] + NOT NULL + +classpath + [VARCHAR(1000)] + NOT NULL + +created_date + [TIMESTAMP] + NOT NULL + +kwargs + [JSON] + NOT NULL + +triggerer_id + [INTEGER] trigger--task_instance - -0..N -{0,1} + +0..N +{0,1} serialized_dag - -serialized_dag + +serialized_dag + +dag_id + [VARCHAR(250)] + NOT NULL -dag_id - [VARCHAR(250)] - NOT NULL +dag_hash + [VARCHAR(32)] + NOT NULL -dag_hash - [VARCHAR(32)] - NOT NULL +data + [JSON] -data - [JSON] +data_compressed + [BLOB] -data_compressed - [BLOB] +fileloc + [VARCHAR(2000)] + NOT NULL -fileloc - [VARCHAR(2000)] - NOT NULL +fileloc_hash + [BIGINT] + NOT NULL -fileloc_hash - [BIGINT] - NOT NULL +last_updated + [TIMESTAMP] + NOT NULL -last_updated - [TIMESTAMP] - NOT NULL - -processor_subdir - [VARCHAR(2000)] +processor_subdir + [VARCHAR(2000)] session - -session + +session + +id + [INTEGER] + NOT NULL -id - [INTEGER] - NOT NULL +data + [BLOB] -data - [BLOB] +expiry + [DATETIME] -expiry - [DATETIME] - -session_id - [VARCHAR(255)] +session_id + [VARCHAR(255)] sla_miss - -sla_miss + +sla_miss + +dag_id + [VARCHAR(250)] + NOT NULL -dag_id - [VARCHAR(250)] - NOT NULL +execution_date + [TIMESTAMP] + NOT NULL -execution_date - [TIMESTAMP] - NOT NULL +task_id + [VARCHAR(250)] + NOT NULL -task_id - [VARCHAR(250)] - NOT NULL +description + [TEXT] -description - [TEXT] +email_sent + [BOOLEAN] -email_sent - [BOOLEAN] +notification_sent + [BOOLEAN] -notification_sent - [BOOLEAN] - -timestamp - [TIMESTAMP] +timestamp + [TIMESTAMP] slot_pool - -slot_pool + +slot_pool + +id + [INTEGER] + NOT NULL -id - [INTEGER] - NOT NULL +description + [TEXT] -description - [TEXT] +pool + [VARCHAR(256)] -pool - [VARCHAR(256)] - -slots - [INTEGER] +slots + [INTEGER] variable - -variable + +variable + +id + [INTEGER] + NOT NULL -id - [INTEGER] - NOT NULL +description + [TEXT] -description - [TEXT] +is_encrypted + [BOOLEAN] -is_encrypted - [BOOLEAN] +key + [VARCHAR(250)] -key - [VARCHAR(250)] - -val - [TEXT] +val + [TEXT] diff --git a/docs/apache-airflow/migrations-ref.rst b/docs/apache-airflow/migrations-ref.rst index 3b92588ba0a90..41b5f5ff04126 100644 --- a/docs/apache-airflow/migrations-ref.rst +++ b/docs/apache-airflow/migrations-ref.rst @@ -39,7 +39,9 @@ Here's the list of all the Database Migrations that are executed via when you ru +---------------------------------+-------------------+-------------------+--------------------------------------------------------------+ | Revision ID | Revises ID | Airflow Version | Description | +=================================+===================+===================+==============================================================+ -| ``ee8d93fcc81e`` (head) | ``ecb43d2a1842`` | ``2.5.0`` | Add updated_at column to DagRun and TaskInstance | +| ``ee8d93fcc81e`` (head) | ``2b72b0fd20ef`` | ``2.5.0`` | Add updated_at column to DagRun and TaskInstance | ++---------------------------------+-------------------+-------------------+--------------------------------------------------------------+ +| ``2b72b0fd20ef`` | ``ecb43d2a1842`` | ``2.4.2`` | Add event timestamp to DDRQ | +---------------------------------+-------------------+-------------------+--------------------------------------------------------------+ | ``ecb43d2a1842`` | ``1486deb605b4`` | ``2.4.0`` | Add processor_subdir column to DagModel, SerializedDagModel | | | | | and CallbackRequest tables. | diff --git a/tests/datasets/test_manager.py b/tests/datasets/test_manager.py index e462d0093e06c..381b17d6009f1 100644 --- a/tests/datasets/test_manager.py +++ b/tests/datasets/test_manager.py @@ -71,16 +71,22 @@ def test_register_dataset_change(self, session, dag_maker, mock_task_instance): dag2 = DagModel(dag_id="dag2") session.add_all([dag1, dag2]) - dsm = DatasetModel(uri="test_dataset_uri") - session.add(dsm) - dsm.consuming_dags = [DagScheduleDatasetReference(dag_id=dag.dag_id) for dag in (dag1, dag2)] + dataset_model = DatasetModel(uri="test_dataset_uri") + session.add(dataset_model) + dataset_model.consuming_dags = [ + DagScheduleDatasetReference(dag_id=dag.dag_id) for dag in (dag1, dag2) + ] session.flush() dsem.register_dataset_change(task_instance=mock_task_instance, dataset=ds, session=session) - # Ensure we've created a dataset - assert session.query(DatasetEvent).filter_by(dataset_id=dsm.id).count() == 1 - assert session.query(DatasetDagRunQueue).count() == 2 + # Ensure we've created a dataset event + events = session.query(DatasetEvent).filter_by(dataset_id=dataset_model.id).all() + assert len(events) == 1 + event = events[0] + queue_records = session.query(DatasetDagRunQueue).all() + assert len(queue_records) == 2 + assert all(x.event_timestamp == event.timestamp for x in queue_records) def test_register_dataset_change_no_downstreams(self, session, mock_task_instance): dsem = DatasetManager() diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 81fd6a3017797..f744d9626e954 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3131,12 +3131,16 @@ def test_create_dag_runs_datasets(self, session, dag_maker): with dag_maker(dag_id="datasets-consumer-single", schedule=[dataset1]): pass dag3 = dag_maker.dag - + session.flush() session = dag_maker.session session.add_all( [ - DatasetDagRunQueue(dataset_id=ds1_id, target_dag_id=dag2.dag_id), - DatasetDagRunQueue(dataset_id=ds1_id, target_dag_id=dag3.dag_id), + DatasetDagRunQueue( + dataset_id=ds1_id, target_dag_id=dag2.dag_id, event_timestamp=event2.timestamp + ), + DatasetDagRunQueue( + dataset_id=ds1_id, target_dag_id=dag3.dag_id, event_timestamp=event2.timestamp + ), ] ) session.flush() diff --git a/tests/models/test_dag.py b/tests/models/test_dag.py index e8f789eafaab8..a77b4da47b20e 100644 --- a/tests/models/test_dag.py +++ b/tests/models/test_dag.py @@ -2234,7 +2234,13 @@ def test_dags_needing_dagruns_datasets(self, dag_maker, session): # add queue records so we'll need a run dag_model = session.query(DagModel).filter(DagModel.dag_id == dag.dag_id).one() dataset_model: DatasetModel = dag_model.schedule_datasets[0] - session.add(DatasetDagRunQueue(dataset_id=dataset_model.id, target_dag_id=dag_model.dag_id)) + session.add( + DatasetDagRunQueue( + dataset_id=dataset_model.id, + target_dag_id=dag_model.dag_id, + event_timestamp=pendulum.now('UTC'), + ) + ) session.flush() query, _ = DagModel.dags_needing_dagruns(session) dag_models = query.all() @@ -2381,11 +2387,20 @@ def test_dags_needing_dagruns_dataset_triggered_dag_info_queued_times(self, sess pass session.flush() + event_timestamp = pendulum.now('UTC') session.add_all( [ - DatasetDagRunQueue(dataset_id=ds1_id, target_dag_id=dag.dag_id, created_at=DEFAULT_DATE), DatasetDagRunQueue( - dataset_id=ds2_id, target_dag_id=dag.dag_id, created_at=DEFAULT_DATE + timedelta(hours=1) + dataset_id=ds1_id, + target_dag_id=dag.dag_id, + created_at=DEFAULT_DATE, + event_timestamp=event_timestamp, + ), + DatasetDagRunQueue( + dataset_id=ds2_id, + target_dag_id=dag.dag_id, + created_at=DEFAULT_DATE + timedelta(hours=1), + event_timestamp=event_timestamp, ), ] ) @@ -2968,10 +2983,11 @@ def test_get_dataset_triggered_next_run_info(dag_maker): session = dag_maker.session ds1_id = session.query(DatasetModel.id).filter_by(uri=dataset1.uri).scalar() + event_timestamp = pendulum.now('UTC') session.bulk_save_objects( [ - DatasetDagRunQueue(dataset_id=ds1_id, target_dag_id=dag2.dag_id), - DatasetDagRunQueue(dataset_id=ds1_id, target_dag_id=dag3.dag_id), + DatasetDagRunQueue(dataset_id=ds1_id, target_dag_id=dag2.dag_id, event_timestamp=event_timestamp), + DatasetDagRunQueue(dataset_id=ds1_id, target_dag_id=dag3.dag_id, event_timestamp=event_timestamp), ] ) session.flush() diff --git a/tests/www/views/test_views_grid.py b/tests/www/views/test_views_grid.py index 9e79319e8b28d..bc28a6ff7232c 100644 --- a/tests/www/views/test_views_grid.py +++ b/tests/www/views/test_views_grid.py @@ -343,7 +343,10 @@ def test_next_run_datasets(admin_client, dag_maker, session, app, monkeypatch): ds1_id = session.query(DatasetModel.id).filter_by(uri=datasets[0].uri).scalar() ds2_id = session.query(DatasetModel.id).filter_by(uri=datasets[1].uri).scalar() ddrq = DatasetDagRunQueue( - target_dag_id=DAG_ID, dataset_id=ds1_id, created_at=pendulum.DateTime(2022, 8, 2, tzinfo=UTC) + target_dag_id=DAG_ID, + dataset_id=ds1_id, + created_at=pendulum.DateTime(2022, 8, 2, tzinfo=UTC), + event_timestamp=pendulum.DateTime(2022, 8, 2, 2, tzinfo=UTC), ) session.add(ddrq) dataset_events = [