Repository navigation
Add AssetAndTimeSchedule timetable - #58543
Conversation
c752a20 to
760844c
Compare
AssetAndTimeScheduleAssetAndTimeSchedule
|
We can probably do some refactoring to extract some common logic between the And and Or timetables into a mixin. This can be a different PR. |
|
Hi @uranusjr, I'd appreciate it if you could take another look when you have time. |
There was a problem hiding this comment.
AI Review
LGTM — the asset-gated scheduling logic (no-placeholder-DagRun retry, ADRQ row-locked re-check, Task SDK flag mirroring without DB leakage, migration symmetry) is correct and well covered by tests across Core and the Task SDK. CI is green. Two minor nits inline, neither blocking.
Smaller observations
airflow-core/src/airflow/jobs/scheduler_job_runner.py:2832— the new shared_lock_queued_asset_recordshelper unconditionally adds.options(joinedload(AssetDagRunQueue.asset)). The pre-existing_create_dag_runs_asset_triggeredcall site (which used to inline this query without the join) 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.airflow-core/tests/unit/timetables/test_assets_timetable.py:398,421-422,450-451—test_infer_manual_data_interval_and,test_next_dagrun_info_and, andtest_generate_run_id_anduseDateTime.now()instead oftime_machine(this repo's testing standard: "Usetime_machinefor time-dependent tests. Do not usedatetime.now()"). This mirrors the existing siblingAssetOrTimeScheduletests 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.
This review was drafted by an AI-assisted tool and confirmed by an Airflow maintainer. The maintainer approving this PR has read the findings and signed off. If something feels off, please reply on the PR and a maintainer will follow up.
More on how Airflow handles maintainer review:
contributing-docs/05_pull_requests.rst.
|
I will review it early tomorrow. Thanks @ashb for the reminder and @aaron-y-chen for the ongoing effort! |
Lee-W
left a comment
There was a problem hiding this comment.
a few nits. nothing major. I'm good with merging it.
|
also need to resolve conflcit |
Add a new built-in timetable that combines a time-based schedule with an asset condition: a SCHEDULED DagRun is created only when both the timetable's scheduled time has arrived and every required asset has queued an event. When the run is created, those asset events are consumed so the next scheduled run waits for fresh updates. Unlike AssetOrTimeSchedule, this does not create asset-triggered runs. Implementation outline: - airflow.timetables.assets.AssetAndTimeSchedule (core) and airflow.sdk.definitions.timetables.assets.AssetAndTimeSchedule (Task SDK shim) delegate next_dagrun_info / infer_manual_data_interval / generate_run_id to the wrapped timetable, retaining the asset_condition that gates creation. - DagModel.dags_needing_dagruns returns a third bucket asset_gated_ready_dag_ids of dags whose asset condition is already satisfied; the SQL query excludes AssetAndTimeSchedule dags whose assets are not yet ready, preventing them from saturating the max_dagruns_to_create_per_loop batch and starving pure time dags. - _create_dagruns_for_dags dispatches the new bucket to _create_dag_runs_asset_gated. The method re-locks ADRQ rows with skip_locked, re-evaluates the asset condition as a HA race guard, creates a SCHEDULED DagRun, links the consumed AssetEvent rows via dag_run.consumed_asset_events.extend(...), updates active_runs and exceeds_max_non_backfill so a follow-up asset event cannot bypass max_active_runs, and finally deletes only the consumed ADRQ rows. - _start_queued_dagruns is untouched: asset gating happens at DagRun creation time, not after. There is no placeholder QUEUED DagRun waiting for assets and no dagrun_timeout-based FAILED path, both of which would have violated the DagRun state-model invariants documented in airflow-core/docs/core-concepts/dag-run.rst. Tests: - AssetAndTimeSchedule serialization/deserialization, data-interval inference, next_dagrun_info delegation, and generate_run_id. - Scheduler-level: time-not-ready, assets-not-ready, both-ready, oldest-pending-slot semantics on late asset arrival, only-consumes-one-slot when multiple schedule slots are pending, ADRQ row locking race guard, max_active_runs respected after a queued run is created, AssetEvent rows linked via consumed_asset_events, and starvation prevention showing the new SQL filter keeps gated dags from starving pure time dags when assets are missing. - DagModel.dags_needing_dagruns regression coverage for the new third-bucket return value. closes: apache#58056
Document the new gating model: a scheduled DagRun is created only when both the timetable's scheduled time has arrived and every required asset has queued an event. If the scheduled time arrives before assets are ready, no DagRun is created; the scheduler holds the oldest pending slot and re-checks each loop until assets arrive. closes: apache#58056
Review feedback asked for the scheduler to handle timetables through behavioral contracts instead of special-cased code paths. The dedicated asset-gated creation path duplicated both the scheduled-run creation logic and the asset-triggered ADRQ/event mechanics, and could drift from either. Asset gating is now a generic step of the normal scheduled-run path, driven only by DagModel.timetable_asset_gated and Timetable.asset_condition, with the ADRQ locking, event-provenance and consumption helpers shared with the asset-triggered path.
…ing validation checks
The previous tests could pass when schedule calculation or run ID delegation returned the wrong value, and wall-clock inputs made the results harder to reason about.
…for improved asset loading
Co-authored-by: Wei Lee <hello@wei-lee.me>
Co-authored-by: Wei Lee <hello@wei-lee.me>
The validation error now starts with an uppercase letter, while the expected regex still used lowercase. This mismatch caused all four nested timetable cases to fail in CI.
Closes: #58056
Why
Use case: #58056
How
1. Add
AssetAndTimeScheduleAdd a timetable that combines a time-based timetable with an asset condition.
It behaves like a scheduled timetable, not an asset-triggered timetable: when it creates a run, the run type is still
SCHEDULED.2. Gate DagRun creation before creating the DagRun
The scheduler creates a scheduled DagRun only when both conditions are satisfied: the schedule time is due and the required asset condition is ready.
If assets are missing, no DagRun is created and
next_dagrun_create_afterstays on the pending slot, so the same slot can be retried later without creating placeholder DagRuns.3. Keep asset-aware scheduler behavior consistent
Timetables expose asset-aware scheduling through
asset_triggeredandasset_gatedbehavior flags, following the existingperiodicpattern. These flags are also mirrored on the Task SDK base timetable so custom timetables can opt in without scheduler type checks.DagModel.dags_needing_dagrunsuses these flags to route asset-triggered Dags separately from asset-gated scheduled Dags. Asset-gated Dags still flow through the standard scheduled-run creation path (_create_dag_runs).When
DagModel.timetable_asset_gatedis set, the scheduler re-checks the timetable'sasset_conditionunder ADRQ row locks before creating the scheduled DagRun. If the condition is still satisfied, the created run links and consumes the selected asset events. The ADRQ locking, event-provenance, and consumption helpers are shared with the asset-triggered path.What
Here are my test DAGs:
upstream_asset_producerdownstream_asset_and_time_consumerFrom the screenshots, the downstream DAG runs as expected: it runs only after both the timetable and the assets are ready.
Upstream dag
Asset production time
Downstream,
AssetAndTimeScheduleImportant
🛠️ Maintainer triage note for @nailo2c · by
@potiuk· 2026-06-22 06:31 UTCYour review threads from
@jscheffllook addressed — please confirm this PR is ready for maintainer review confirmation:The ball is in your court — you've been assigned to this PR. Reply
yes / ready(and mark the threads resolved) and a maintainer will pick it up from the queue.Automated triage — may be imperfect; a maintainer takes the next look.