kaxil commented on code in PR #73952:
URL: https://github.com/apache/airflow/pull/73952#discussion_r4189871192
##########
airflow-core/src/airflow/models/dag.py:
##########
@@ -806,6 +811,9 @@ def dag_ready(dag_id: str, cond: SerializedAssetBase,
statuses: dict[UKey, bool]
.order_by(cls.next_dagrun_create_after)
.limit(cls.NUM_DAGS_PER_DAGRUN_QUERY)
)
+ if not include_scheduled:
+ # Dags due by their timetable get no run, so they would stay due
and could fill the batch.
Review Comment:
This filter is now the only thing enforcing `use_job_schedule=False`.
`_create_dagruns_for_dags` hands every non-asset Dag to `_create_dag_runs`,
which doesn't check the flag, so without it timetable-due Dags would get
scheduled runs. The comment reads as if they get no run anyway and this is only
about filling the batch, which makes it look safe to drop. Could the comment
say that? Building the predicate once before the query (just
`cls.dag_id.in_(asset_triggered_dag_ids)` when `include_scheduled` is False,
the `or_` otherwise) would also avoid a `.where()` after `.limit()` that reads
like a post-limit filter, even though SQLAlchemy puts it in the same WHERE.
##########
airflow-core/docs/best-practices.rst:
##########
@@ -933,7 +933,7 @@ Disable the scheduler
You might consider disabling the Airflow cluster while you perform such
maintenance.
-One way to do so would be to set the param ``[scheduler] > use_job_schedule``
to ``False`` and wait for any running Dags to complete; after this no new Dag
runs will be created unless externally triggered.
+One way to do so would be to set the param ``[scheduler] > use_job_schedule``
to ``False`` and wait for any running Dags to complete; after this no new Dag
runs will be created on schedule, though manually triggered and asset-triggered
Dag runs are still created.
Review Comment:
The `use_job_schedule` description in `config.yml` (what users see in the
config reference) still says only manually triggered Dags keep running. Could
you update it with this tip? Something like "Set to ``False`` to stop the
scheduler creating Dag runs from Dag timetables. Manually triggered and
asset-triggered Dag runs are still created." "Cron intervals" was also narrower
than what the flag covers, since timedelta and other timetables stop too.
##########
airflow-core/tests/unit/models/test_dag.py:
##########
@@ -2690,6 +2690,42 @@ def test_dags_needing_dagruns_assets(self, dag_maker,
session):
dag_models = query.all()
assert dag_models == [dag_model]
+ def test_dags_needing_dagruns_without_scheduled(self, dag_maker, session):
+ """include_scheduled=False leaves out a Dag due by its schedule but
keeps an asset-triggered one."""
+ asset = Asset(uri="test://asset-without-scheduled", group="test-group")
+ with dag_maker(
+ session=session,
+ dag_id="due_by_schedule",
+ schedule="@daily",
+ start_date=pendulum.now().add(days=-2),
+ ):
+ EmptyOperator(task_id="dummy")
+ assert dag_maker.dag_model.next_dagrun_create_after <=
timezone.utcnow()
+ with dag_maker(
+ session=session,
+ dag_id="asset_triggered",
+ schedule=[asset],
+ start_date=pendulum.now().add(days=-2),
+ ):
+ EmptyOperator(task_id="dummy")
+ asset_model = dag_maker.dag_model.schedule_assets[0]
+ event = AssetEvent(asset_id=asset_model.id,
timestamp=timezone.utcnow())
+ session.add(event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_model.id, target_dag_id="asset_triggered",
asset_event_id=event.id
+ )
+ )
+ session.flush()
+
+ query, _ = DagModel.dags_needing_dagruns(session)
+ assert sorted(dag_model.dag_id for dag_model in query) ==
["asset_triggered", "due_by_schedule"]
+
+ query, triggered_date_by_dag = DagModel.dags_needing_dagruns(session,
include_scheduled=False)
Review Comment:
With two Dags and the default limit of 10, a filter applied in Python after
the query would pass this too, so it doesn't pin the batch behaviour the
description mentions. Patching `DagModel.NUM_DAGS_PER_DAGRUN_QUERY` to 1 for
this call would: on Postgres the asset Dag's NULL `next_dagrun_create_after`
sorts last, so the due Dag takes the only slot unless the filter is in SQL.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]