shahar1 commented on code in PR #58543:
URL: https://github.com/apache/airflow/pull/58543#discussion_r4102259952
##########
airflow-core/tests/unit/timetables/test_assets_timetable.py:
##########
@@ -197,6 +394,12 @@ def test_infer_manual_data_interval(core_asset_timetable:
CoreAssetOrTimeSchedul
assert result == DataInterval.exact(run_after)
+def test_infer_manual_data_interval_and(core_asset_and_time_timetable:
CoreAssetAndTimeSchedule) -> None:
+ run_after = DateTime.now()
+ result =
core_asset_and_time_timetable.infer_manual_data_interval(run_after=run_after)
Review Comment:
`test_infer_manual_data_interval_and` (line 398),
`test_next_dagrun_info_and` (lines 421-422), and `test_generate_run_id_and`
(lines 450-451) use `DateTime.now()` instead of `time_machine`, against this
repo's testing standard ("Use `time_machine` for time-dependent tests. Do not
use `datetime.now()`"). Mirrors the existing sibling `AssetOrTimeSchedule`
tests in the same file rather than introducing a new deviation, and the
assertions are structural (`isinstance(...)`) so it isn't flaky today — but new
code shouldn't reproduce the pattern.
##########
airflow-core/src/airflow/jobs/scheduler_job_runner.py:
##########
@@ -2776,6 +2793,113 @@ def _create_dag_runs(self, dag_models:
Collection[DagModel], session: Session) -
# TODO[HA]: Should we do a session.flush() so we don't have to
keep lots of state/object in
# memory for larger dags? or expunge_all()
+ def _collect_gated_asset_events(
+ self, *, dag: SerializedDAG, session: Session
+ ) -> tuple[Sequence[AssetDagRunQueue], list[AssetEvent]] | None:
+ """
+ Check an asset-gated Dag's asset condition and collect what a new run
consumes.
+
+ Returns ``None`` when the condition is not satisfied by the queued
asset
+ events, in which case no run should be created yet.
+ """
+ records = self._lock_queued_asset_records(dag_id=dag.dag_id,
session=session)
+ if not records:
+ return None
+ statuses = {SerializedAssetUniqueKey.from_asset(record.asset): True
for record in records}
+ try:
+ ready = AssetEvaluator(session).run(dag.timetable.asset_condition,
statuses=statuses)
+ except Exception:
+ self.log.exception("Dag '%s' failed to be evaluated; assuming not
ready", dag.dag_id)
+ return None
+ if not ready:
+ return None
+ asset_events = self._select_consumed_asset_events(
+ dag=dag,
+ records=records,
+ session=session,
+ )
+ if not asset_events:
+ self._delete_consumed_asset_records(records=records,
dag_id=dag.dag_id, session=session)
+ return None
+ return records, asset_events
+
+ def _lock_queued_asset_records(self, *, dag_id: str, session: Session) ->
Sequence[AssetDagRunQueue]:
+ """Lock and return the Dag's queued asset (ADRQ) rows, skipping rows
another scheduler holds."""
+ return session.scalars(
+ with_row_locks(
+ select(AssetDagRunQueue)
+ .where(AssetDagRunQueue.target_dag_id == dag_id)
+ .options(joinedload(AssetDagRunQueue.asset)),
+ of=AssetDagRunQueue,
Review Comment:
Shared `_lock_queued_asset_records` now unconditionally adds
`.options(joinedload(AssetDagRunQueue.asset))`. The pre-existing
`_create_dag_runs_asset_triggered` call site (which used to inline this query
without the join, before this refactor) never reads `.asset`, so it now pays
for an extra JOIN/hydration on every asset-triggered scheduler pass for no
benefit. Not a correctness issue, just a small efficiency regression from the
refactor — worth gating the eager load behind a parameter if it matters at
scale.
--
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]