From beb01e4ca0d1ed1e38087140be0e0614cb2cddfe Mon Sep 17 00:00:00 2001 From: tosheer Date: Tue, 14 May 2024 10:25:37 +0530 Subject: [PATCH 1/9] Fixed issue of new dag getting old dataset events. --- airflow/jobs/scheduler_job_runner.py | 2 + tests/jobs/test_scheduler_job.py | 69 +++++++++++++++++++++------- 2 files changed, 55 insertions(+), 16 deletions(-) diff --git a/airflow/jobs/scheduler_job_runner.py b/airflow/jobs/scheduler_job_runner.py index f2333e8d5a0d6..dd730c9855efb 100644 --- a/airflow/jobs/scheduler_job_runner.py +++ b/airflow/jobs/scheduler_job_runner.py @@ -1277,6 +1277,8 @@ def _create_dag_runs_dataset_triggered( ] if previous_dag_run: dataset_event_filters.append(DatasetEvent.timestamp > previous_dag_run.execution_date) + else: + dataset_event_filters.append(DatasetEvent.timestamp > DagScheduleDatasetReference.created_at) dataset_events = session.scalars( select(DatasetEvent) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 85399892acdc0..8f2d7bb227836 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3776,6 +3776,11 @@ def test_create_dag_runs_datasets(self, session, dag_maker): dataset1 = Dataset(uri="ds1") dataset2 = Dataset(uri="ds2") + # Create DAG before the arrival of dataset events. + with dag_maker(dag_id="datasets-consumer-single-old", schedule=[dataset1]): + pass + dag2 = dag_maker.dag + with dag_maker(dag_id="datasets-1", start_date=timezone.utcnow(), session=session): BashOperator(task_id="task", bash_command="echo 1", outlets=[dataset1]) dr = dag_maker.create_dagrun( @@ -3811,20 +3816,40 @@ def test_create_dag_runs_datasets(self, session, dag_maker): ) session.add(event2) + # Create a third event, creation time is more recent, but data interval is even older + dr = dag_maker.create_dagrun( + run_id="run3", + execution_date=(DEFAULT_DATE + timedelta(days=102)), + data_interval=(DEFAULT_DATE + timedelta(days=3), DEFAULT_DATE + timedelta(days=3)), + ) + + event3 = DatasetEvent( + dataset_id=ds1_id, + source_task_id="task", + source_dag_id=dr.dag_id, + source_run_id=dr.run_id, + source_map_index=-1, + ) + session.add(event3) + with dag_maker(dag_id="datasets-consumer-multiple", schedule=[dataset1, dataset2]): pass - dag2 = dag_maker.dag - with dag_maker(dag_id="datasets-consumer-single", schedule=[dataset1]): - pass dag3 = dag_maker.dag + # Create DAG after dataset events. + with dag_maker(dag_id="datasets-consumer-single-new", schedule=[dataset1]): + pass + dag4 = dag_maker.dag + 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=dag4.dag_id), ] ) + session.flush() scheduler_job = Job(executor=self.null_exec) @@ -3838,27 +3863,39 @@ def test_create_dag_runs_datasets(self, session, dag_maker): def dict_from_obj(obj): """Get dict of column attrs from SqlAlchemy object.""" return {k.key: obj.__dict__.get(k) for k in obj.__mapper__.column_attrs} - - # dag3 should be triggered since it only depends on dataset1, and it's been queued - created_run = session.query(DagRun).filter(DagRun.dag_id == dag3.dag_id).one() + + # dag2 should be triggered since it only depends on dataset1, it's been queued and dataset events landed after DAG was created. + created_run = session.query(DagRun).filter(DagRun.dag_id == dag2.dag_id).one() assert created_run.state == State.QUEUED assert created_run.start_date is None # we don't have __eq__ defined on DatasetEvent because... given the fact that in the future # we may register events from other systems, dataset_id + timestamp might not be enough PK assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == list( - map(dict_from_obj, [event1, event2]) + map(dict_from_obj, [event1, event2, event3]) ) - assert created_run.data_interval_start == DEFAULT_DATE + timedelta(days=5) + assert created_run.data_interval_start == DEFAULT_DATE + timedelta(days=3) assert created_run.data_interval_end == DEFAULT_DATE + timedelta(days=11) - # dag2 DDRQ record should still be there since the dag run was *not* triggered - assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag2.dag_id).one() is not None - # dag2 should not be triggered since it depends on both dataset 1 and 2 - assert session.query(DagRun).filter(DagRun.dag_id == dag2.dag_id).one_or_none() is None - # dag3 DDRQ record should be deleted since the dag run was triggered - assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag3.dag_id).one_or_none() is None - - assert dag3.get_last_dagrun().creating_job_id == scheduler_job.id + + # dag4 should be triggered since it only depends on dataset1, and it's been queued + created_run_dag4 = session.query(DagRun).filter(DagRun.dag_id == dag4.dag_id).one() + assert created_run_dag4.state == State.QUEUED + assert created_run_dag4.start_date is None + + # we don't have __eq__ defined on DatasetEvent because... given the fact that in the future + # we may register events from other systems, dataset_id + timestamp might not be enough PK + assert list(map(dict_from_obj, created_run_dag4.consumed_dataset_events)) == list( + map(dict_from_obj, []) + ) + + # dag3 DDRQ record should still be there since the dag run was *not* triggered + assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag3.dag_id).one() is not None + # dag3 should not be triggered since it depends on both dataset 1 and 2 + assert session.query(DagRun).filter(DagRun.dag_id == dag3.dag_id).one_or_none() is None + # dag2 DDRQ record should be deleted since the dag run was triggered + assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag2.dag_id).one_or_none() is None + + assert dag2.get_last_dagrun().creating_job_id == scheduler_job.id @pytest.mark.need_serialized_dag @pytest.mark.parametrize( From 1f8cb5238d9a01d303b415d3779af2aee7a32cd4 Mon Sep 17 00:00:00 2001 From: tosheer Date: Tue, 14 May 2024 16:11:52 +0530 Subject: [PATCH 2/9] fixed static checks. --- airflow/jobs/scheduler_job_runner.py | 4 +++- tests/jobs/test_scheduler_job.py | 2 +- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/airflow/jobs/scheduler_job_runner.py b/airflow/jobs/scheduler_job_runner.py index dd730c9855efb..51c6d485ce85d 100644 --- a/airflow/jobs/scheduler_job_runner.py +++ b/airflow/jobs/scheduler_job_runner.py @@ -1278,7 +1278,9 @@ def _create_dag_runs_dataset_triggered( if previous_dag_run: dataset_event_filters.append(DatasetEvent.timestamp > previous_dag_run.execution_date) else: - dataset_event_filters.append(DatasetEvent.timestamp > DagScheduleDatasetReference.created_at) + dataset_event_filters.append( + DatasetEvent.timestamp > DagScheduleDatasetReference.created_at + ) dataset_events = session.scalars( select(DatasetEvent) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 8f2d7bb227836..b5b2de379ef7a 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3863,7 +3863,7 @@ def test_create_dag_runs_datasets(self, session, dag_maker): def dict_from_obj(obj): """Get dict of column attrs from SqlAlchemy object.""" return {k.key: obj.__dict__.get(k) for k in obj.__mapper__.column_attrs} - + # dag2 should be triggered since it only depends on dataset1, it's been queued and dataset events landed after DAG was created. created_run = session.query(DagRun).filter(DagRun.dag_id == dag2.dag_id).one() assert created_run.state == State.QUEUED From 4346f3acc5c709fc00dae5e2096afa2c04cc06c7 Mon Sep 17 00:00:00 2001 From: tosheer Date: Thu, 16 May 2024 11:49:03 +0530 Subject: [PATCH 3/9] added additional test case. --- airflow/jobs/scheduler_job_runner.py | 2 +- tests/jobs/test_scheduler_job.py | 149 ++++++++++++++++++++------- 2 files changed, 110 insertions(+), 41 deletions(-) diff --git a/airflow/jobs/scheduler_job_runner.py b/airflow/jobs/scheduler_job_runner.py index 51c6d485ce85d..d468f4f693a45 100644 --- a/airflow/jobs/scheduler_job_runner.py +++ b/airflow/jobs/scheduler_job_runner.py @@ -1279,7 +1279,7 @@ def _create_dag_runs_dataset_triggered( dataset_event_filters.append(DatasetEvent.timestamp > previous_dag_run.execution_date) else: dataset_event_filters.append( - DatasetEvent.timestamp > DagScheduleDatasetReference.created_at + DatasetEvent.timestamp >= DagScheduleDatasetReference.created_at ) dataset_events = session.scalars( diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index b5b2de379ef7a..66891029fa933 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3776,8 +3776,7 @@ def test_create_dag_runs_datasets(self, session, dag_maker): dataset1 = Dataset(uri="ds1") dataset2 = Dataset(uri="ds2") - # Create DAG before the arrival of dataset events. - with dag_maker(dag_id="datasets-consumer-single-old", schedule=[dataset1]): + with dag_maker(dag_id="datasets-consumer-single", schedule=[dataset1]): pass dag2 = dag_maker.dag @@ -3816,40 +3815,17 @@ def test_create_dag_runs_datasets(self, session, dag_maker): ) session.add(event2) - # Create a third event, creation time is more recent, but data interval is even older - dr = dag_maker.create_dagrun( - run_id="run3", - execution_date=(DEFAULT_DATE + timedelta(days=102)), - data_interval=(DEFAULT_DATE + timedelta(days=3), DEFAULT_DATE + timedelta(days=3)), - ) - - event3 = DatasetEvent( - dataset_id=ds1_id, - source_task_id="task", - source_dag_id=dr.dag_id, - source_run_id=dr.run_id, - source_map_index=-1, - ) - session.add(event3) - with dag_maker(dag_id="datasets-consumer-multiple", schedule=[dataset1, dataset2]): pass dag3 = dag_maker.dag - # Create DAG after dataset events. - with dag_maker(dag_id="datasets-consumer-single-new", schedule=[dataset1]): - pass - dag4 = dag_maker.dag - 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=dag4.dag_id), ] ) - session.flush() scheduler_job = Job(executor=self.null_exec) @@ -3872,30 +3848,123 @@ def dict_from_obj(obj): # we don't have __eq__ defined on DatasetEvent because... given the fact that in the future # we may register events from other systems, dataset_id + timestamp might not be enough PK assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == list( - map(dict_from_obj, [event1, event2, event3]) + map(dict_from_obj, [event1, event2]) ) - assert created_run.data_interval_start == DEFAULT_DATE + timedelta(days=3) + assert created_run.data_interval_start == DEFAULT_DATE + timedelta(days=5) assert created_run.data_interval_end == DEFAULT_DATE + timedelta(days=11) + # dag2 DDRQ record should still be there since the dag run was *not* triggered + assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag3.dag_id).one() is not None + # dag2 should not be triggered since it depends on both dataset 1 and 2 + assert session.query(DagRun).filter(DagRun.dag_id == dag3.dag_id).one_or_none() is None + # dag3 DDRQ record should be deleted since the dag run was triggered + assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag2.dag_id).one_or_none() is None - # dag4 should be triggered since it only depends on dataset1, and it's been queued - created_run_dag4 = session.query(DagRun).filter(DagRun.dag_id == dag4.dag_id).one() - assert created_run_dag4.state == State.QUEUED - assert created_run_dag4.start_date is None + assert dag2.get_last_dagrun().creating_job_id == scheduler_job.id + + @pytest.mark.need_serialized_dag + def test_new_dagrun_ignores_old_dataset_events(self, session, dag_maker): + """ + Test various invariants of _create_dag_runs. + + - That the new DAG should not get dataset events which has timestamp with before dag creation date. + - That the run created is on QUEUED State + - That dag_model has next_dagrun + """ + + dataset = Dataset(uri="ds") + + with dag_maker(dag_id="datasets-1", start_date=timezone.utcnow(), session=session): + BashOperator(task_id="task", bash_command="echo 1", outlets=[dataset]) + dr = dag_maker.create_dagrun( + run_id="run1", + execution_date=(DEFAULT_DATE + timedelta(days=100)), + data_interval=(DEFAULT_DATE + timedelta(days=10), DEFAULT_DATE + timedelta(days=11)), + ) + + ds_id = session.query(DatasetModel.id).filter_by(uri=dataset.uri).scalar() + + event1 = DatasetEvent( + dataset_id=ds_id, + source_task_id="task", + source_dag_id=dr.dag_id, + source_run_id=dr.run_id, + source_map_index=-1, + ) + session.add(event1) + + # Create a second event, creation time is more recent, but data interval is older + dr = dag_maker.create_dagrun( + run_id="run2", + execution_date=(DEFAULT_DATE + timedelta(days=101)), + data_interval=(DEFAULT_DATE + timedelta(days=5), DEFAULT_DATE + timedelta(days=6)), + ) + + event2 = DatasetEvent( + dataset_id=ds_id, + source_task_id="task", + source_dag_id=dr.dag_id, + source_run_id=dr.run_id, + source_map_index=-1, + ) + session.add(event2) + + # Create DAG after dataset events. + with dag_maker(dag_id="datasets-consumer", schedule=[dataset]): + pass + dag = dag_maker.dag + + with dag_maker(dag_id="datasets-1-new", start_date=timezone.utcnow(), session=session): + BashOperator(task_id="task", bash_command="echo 1", outlets=[dataset]) + dr = dag_maker.create_dagrun( + run_id="run3", + execution_date=(DEFAULT_DATE + timedelta(days=101)), + data_interval=(DEFAULT_DATE + timedelta(days=5), DEFAULT_DATE + timedelta(days=6)), + ) + + event3 = DatasetEvent( + dataset_id=ds_id, + source_task_id="task", + source_dag_id=dr.dag_id, + source_run_id=dr.run_id, + source_map_index=-1, + timestamp=timezone.utcnow(), + ) + session.add(event3) + + session = dag_maker.session + session.add_all( + [ + DatasetDagRunQueue(dataset_id=ds_id, target_dag_id=dag.dag_id), + ] + ) + session.flush() + + scheduler_job = Job(executor=self.null_exec) + self.job_runner = SchedulerJobRunner(job=scheduler_job) + + self.job_runner.processor_agent = mock.MagicMock() + + with create_session() as session: + self.job_runner._create_dagruns_for_dags(session, session) + + def dict_from_obj(obj): + """Get dict of column attrs from SqlAlchemy object.""" + return {k.key: obj.__dict__.get(k) for k in obj.__mapper__.column_attrs} + + # dag should be triggered since it only depends on dataset1, it's been queued and dataset events landed after DAG was created. + created_run = session.query(DagRun).filter(DagRun.dag_id == dag.dag_id).one() + assert created_run.state == State.QUEUED # we don't have __eq__ defined on DatasetEvent because... given the fact that in the future # we may register events from other systems, dataset_id + timestamp might not be enough PK - assert list(map(dict_from_obj, created_run_dag4.consumed_dataset_events)) == list( - map(dict_from_obj, []) + assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == list( + map(dict_from_obj, [event3]) ) - # dag3 DDRQ record should still be there since the dag run was *not* triggered - assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag3.dag_id).one() is not None - # dag3 should not be triggered since it depends on both dataset 1 and 2 - assert session.query(DagRun).filter(DagRun.dag_id == dag3.dag_id).one_or_none() is None - # dag2 DDRQ record should be deleted since the dag run was triggered - assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag2.dag_id).one_or_none() is None + # dag DDRQ record should be deleted since the dag run was triggered + assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag.dag_id).one_or_none() is None - assert dag2.get_last_dagrun().creating_job_id == scheduler_job.id + assert dag.get_last_dagrun().creating_job_id == scheduler_job.id @pytest.mark.need_serialized_dag @pytest.mark.parametrize( From 681f22c0481389d1dbfe92eb562bd2130b4bc704 Mon Sep 17 00:00:00 2001 From: tosheer Date: Thu, 16 May 2024 12:49:06 +0530 Subject: [PATCH 4/9] added newsfragment --- newsfragments/39603.significant.rst | 8 ++++++++ 1 file changed, 8 insertions(+) create mode 100644 newsfragments/39603.significant.rst diff --git a/newsfragments/39603.significant.rst b/newsfragments/39603.significant.rst new file mode 100644 index 0000000000000..67f42eef83792 --- /dev/null +++ b/newsfragments/39603.significant.rst @@ -0,0 +1,8 @@ +New DAGs will no longer receive historical dataset events + +In the past, new DAGs received historical dataset events from the +very first dataset event in the system. This has been corrected, +and now new DAGs will only receive dataset events that have occurred +after their creation. Although this change may disrupt existing workflows, +the previous behavior is considered a bug, and this update ensures more +accurate and expected DAG execution. From e0745a6533553e5a55685a85f7ed11ebd363a4f5 Mon Sep 17 00:00:00 2001 From: tosheer Date: Wed, 22 May 2024 01:11:27 +0530 Subject: [PATCH 5/9] Incorporated feedback and suggestions. --- airflow/jobs/scheduler_job_runner.py | 2 +- tests/jobs/test_scheduler_job.py | 12 +++--------- 2 files changed, 4 insertions(+), 10 deletions(-) diff --git a/airflow/jobs/scheduler_job_runner.py b/airflow/jobs/scheduler_job_runner.py index a1372719b37fa..1c93ccdafc280 100644 --- a/airflow/jobs/scheduler_job_runner.py +++ b/airflow/jobs/scheduler_job_runner.py @@ -1279,7 +1279,7 @@ def _create_dag_runs_dataset_triggered( dataset_event_filters.append(DatasetEvent.timestamp > previous_dag_run.execution_date) else: dataset_event_filters.append( - DatasetEvent.timestamp >= DagScheduleDatasetReference.created_at + DatasetEvent.timestamp > DagScheduleDatasetReference.created_at ) dataset_events = session.scalars( diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index d797370b7e6fb..b433f2c12d615 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3949,11 +3949,7 @@ def test_new_dagrun_ignores_old_dataset_events(self, session, dag_maker): session.add(event3) session = dag_maker.session - session.add_all( - [ - DatasetDagRunQueue(dataset_id=ds_id, target_dag_id=dag.dag_id), - ] - ) + session.add(DatasetDagRunQueue(dataset_id=ds_id, target_dag_id=dag.dag_id)) session.flush() scheduler_job = Job(executor=self.null_exec) @@ -3966,7 +3962,7 @@ def test_new_dagrun_ignores_old_dataset_events(self, session, dag_maker): def dict_from_obj(obj): """Get dict of column attrs from SqlAlchemy object.""" - return {k.key: obj.__dict__.get(k) for k in obj.__mapper__.column_attrs} + return {k.key: getattr(obj, k) for k in obj.__mapper__.column_attrs} # dag should be triggered since it only depends on dataset1, it's been queued and dataset events landed after DAG was created. created_run = session.query(DagRun).filter(DagRun.dag_id == dag.dag_id).one() @@ -3974,9 +3970,7 @@ def dict_from_obj(obj): # we don't have __eq__ defined on DatasetEvent because... given the fact that in the future # we may register events from other systems, dataset_id + timestamp might not be enough PK - assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == list( - map(dict_from_obj, [event3]) - ) + assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == [dict_from_obj(event3)] # dag DDRQ record should be deleted since the dag run was triggered assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag.dag_id).one_or_none() is None From 71c8d8c7b65ae1c0d10cea60411cc240a541905c Mon Sep 17 00:00:00 2001 From: tosheer Date: Tue, 11 Jun 2024 16:00:25 +0530 Subject: [PATCH 6/9] Update tests/jobs/test_scheduler_job.py Co-authored-by: Ryan Hatter <25823361+RNHTTR@users.noreply.github.com> --- tests/jobs/test_scheduler_job.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index b433f2c12d615..9c77ad5c4f3dd 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3968,8 +3968,7 @@ def dict_from_obj(obj): created_run = session.query(DagRun).filter(DagRun.dag_id == dag.dag_id).one() assert created_run.state == State.QUEUED - # we don't have __eq__ defined on DatasetEvent because... given the fact that in the future - # we may register events from other systems, dataset_id + timestamp might not be enough PK + # __eq__ isn't defined on DatasetEvent assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == [dict_from_obj(event3)] # dag DDRQ record should be deleted since the dag run was triggered From 355620015dcffe356bc968a868d73eec6bdfa3b8 Mon Sep 17 00:00:00 2001 From: tosheer Date: Tue, 11 Jun 2024 16:06:46 +0530 Subject: [PATCH 7/9] updated impl. --- tests/jobs/test_scheduler_job.py | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 9c77ad5c4f3dd..7d48ae87c89e0 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3928,7 +3928,7 @@ def test_new_dagrun_ignores_old_dataset_events(self, session, dag_maker): # Create DAG after dataset events. with dag_maker(dag_id="datasets-consumer", schedule=[dataset]): pass - dag = dag_maker.dag + consumer_dag = dag_maker.dag with dag_maker(dag_id="datasets-1-new", start_date=timezone.utcnow(), session=session): BashOperator(task_id="task", bash_command="echo 1", outlets=[dataset]) @@ -3949,7 +3949,7 @@ def test_new_dagrun_ignores_old_dataset_events(self, session, dag_maker): session.add(event3) session = dag_maker.session - session.add(DatasetDagRunQueue(dataset_id=ds_id, target_dag_id=dag.dag_id)) + session.add(DatasetDagRunQueue(dataset_id=ds_id, target_dag_id=consumer_dag.dag_id)) session.flush() scheduler_job = Job(executor=self.null_exec) @@ -3965,16 +3965,16 @@ def dict_from_obj(obj): return {k.key: getattr(obj, k) for k in obj.__mapper__.column_attrs} # dag should be triggered since it only depends on dataset1, it's been queued and dataset events landed after DAG was created. - created_run = session.query(DagRun).filter(DagRun.dag_id == dag.dag_id).one() + created_run = session.query(DagRun).filter(DagRun.dag_id == consumer_dag.dag_id).one() assert created_run.state == State.QUEUED # __eq__ isn't defined on DatasetEvent assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == [dict_from_obj(event3)] # dag DDRQ record should be deleted since the dag run was triggered - assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=dag.dag_id).one_or_none() is None + assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=consumer_dag.dag_id).one_or_none() is None - assert dag.get_last_dagrun().creating_job_id == scheduler_job.id + assert consumer_dag.get_last_dagrun().creating_job_id == scheduler_job.id @pytest.mark.need_serialized_dag @pytest.mark.parametrize( From 519337f05ddb286c983be91c68fd8a5ff6dde9d6 Mon Sep 17 00:00:00 2001 From: tosheer Date: Fri, 14 Jun 2024 13:54:18 +0530 Subject: [PATCH 8/9] updated for static code check. --- tests/jobs/test_scheduler_job.py | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index 7d48ae87c89e0..ebcfe56a6087f 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3972,7 +3972,10 @@ def dict_from_obj(obj): assert list(map(dict_from_obj, created_run.consumed_dataset_events)) == [dict_from_obj(event3)] # dag DDRQ record should be deleted since the dag run was triggered - assert session.query(DatasetDagRunQueue).filter_by(target_dag_id=consumer_dag.dag_id).one_or_none() is None + assert ( + session.query(DatasetDagRunQueue).filter_by(target_dag_id=consumer_dag.dag_id).one_or_none() + is None + ) assert consumer_dag.get_last_dagrun().creating_job_id == scheduler_job.id From 10658742c1b35a5f46c90bb7f4722f56c6b5fb3c Mon Sep 17 00:00:00 2001 From: tosheer Date: Tue, 9 Jul 2024 12:08:52 +0530 Subject: [PATCH 9/9] updated test case. --- tests/jobs/test_scheduler_job.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/jobs/test_scheduler_job.py b/tests/jobs/test_scheduler_job.py index add305fd716bf..9d1bd853eca59 100644 --- a/tests/jobs/test_scheduler_job.py +++ b/tests/jobs/test_scheduler_job.py @@ -3962,7 +3962,7 @@ def test_new_dagrun_ignores_old_dataset_events(self, session, dag_maker): def dict_from_obj(obj): """Get dict of column attrs from SqlAlchemy object.""" - return {k.key: getattr(obj, k) for k in obj.__mapper__.column_attrs} + return {k.key: obj.__dict__.get(k) for k in obj.__mapper__.column_attrs} # dag should be triggered since it only depends on dataset1, it's been queued and dataset events landed after DAG was created. created_run = session.query(DagRun).filter(DagRun.dag_id == consumer_dag.dag_id).one()