This is an automated email from the ASF dual-hosted git repository.
Lee-W pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new e33e45fe67a Add `AssetAndTimeSchedule` timetable (#58543)
e33e45fe67a is described below
commit e33e45fe67aeaabe797a97c83cfb92d15c36e8a5
Author: Aaron Chen <[email protected]>
AuthorDate: Wed Sep 30 19:22:37 2026 -0700
Add `AssetAndTimeSchedule` timetable (#58543)
Co-authored-by: Wei Lee <[email protected]>
---
.../authoring-and-scheduling/asset-scheduling.rst | 31 +-
.../docs/authoring-and-scheduling/timetable.rst | 25 +-
airflow-core/docs/migrations-ref.rst | 4 +-
airflow-core/newsfragments/58543.feature.rst | 1 +
.../src/airflow/dag_processing/collection.py | 1 +
.../src/airflow/jobs/scheduler_job_runner.py | 207 +++++---
..._3_4_0_add_timetable_asset_gated_to_dagmodel.py | 55 ++
airflow-core/src/airflow/models/dag.py | 32 +-
airflow-core/src/airflow/serialization/encoders.py | 9 +
airflow-core/src/airflow/timetables/assets.py | 112 +++-
airflow-core/src/airflow/timetables/base.py | 16 +
airflow-core/src/airflow/timetables/simple.py | 1 +
airflow-core/src/airflow/utils/db.py | 2 +-
.../tests/unit/dag_processing/test_collection.py | 34 +-
airflow-core/tests/unit/jobs/test_scheduler_job.py | 566 ++++++++++++++++++++-
airflow-core/tests/unit/models/test_dag.py | 12 +
.../tests/unit/timetables/test_assets_timetable.py | 227 ++++++++-
task-sdk/docs/api.rst | 2 +
task-sdk/src/airflow/sdk/__init__.py | 3 +
task-sdk/src/airflow/sdk/__init__.pyi | 2 +
task-sdk/src/airflow/sdk/bases/timetable.py | 7 +
task-sdk/src/airflow/sdk/definitions/dag.py | 2 +-
.../airflow/sdk/definitions/timetables/assets.py | 25 +
task-sdk/tests/task_sdk/definitions/test_dag.py | 16 +-
24 files changed, 1292 insertions(+), 100 deletions(-)
diff --git a/airflow-core/docs/authoring-and-scheduling/asset-scheduling.rst
b/airflow-core/docs/authoring-and-scheduling/asset-scheduling.rst
index 1bfe98d715e..1bf2507237b 100644
--- a/airflow-core/docs/authoring-and-scheduling/asset-scheduling.rst
+++ b/airflow-core/docs/authoring-and-scheduling/asset-scheduling.rst
@@ -147,6 +147,30 @@ If one asset is updated multiple times before all consumed
assets update, the do
}
+Gate scheduled runs on asset updates
+------------------------------------
+
+Use ``AssetAndTimeSchedule`` when you want a Dag to follow a normal time-based
timetable but only create a scheduled DagRun after specific assets have been
updated. Airflow creates the scheduled DagRun only when both the timetable's
scheduled time has arrived and every required asset has queued an event. When
the DagRun is created, those asset events are consumed so the next scheduled
run waits for new updates. This does not create additional asset-triggered runs.
+
+If the scheduled time arrives before the required assets are ready, Airflow
does not create a DagRun. The scheduler re-checks on later loops and holds the
oldest pending scheduled slot until the assets arrive.
+
+.. code-block:: python
+
+ from airflow.sdk import DAG, Asset, AssetAndTimeSchedule,
CronTriggerTimetable
+
+ example_asset = Asset("s3://asset/example.csv")
+
+ with DAG(
+ dag_id="gated_hourly_dag",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("0 * * * *", timezone="UTC"),
+ assets=[example_asset],
+ ),
+ ...,
+ ):
+ ...
+
+
Fetching information from a triggering asset event
----------------------------------------------------
@@ -416,6 +440,9 @@ Combining asset and time-based schedules
AssetTimetable Integration
~~~~~~~~~~~~~~~~~~~~~~~~~~~~
-You can schedule Dags based on both asset events and time-based schedules
using ``AssetOrTimeSchedule``. This allows you to create workflows when a Dag
needs both to be triggered by data updates and run periodically according to a
fixed timetable.
+Asset-aware timetables combine asset expressions with a time-based schedule:
+
+* Use ``AssetOrTimeSchedule`` to create runs both on a timetable and when
assets update, producing scheduled runs and asset-triggered runs independently.
+* Use ``AssetAndTimeSchedule`` to keep a Dag on a timetable but only create
scheduled runs once the referenced assets have been updated.
-For more detailed information on ``AssetOrTimeSchedule``, refer to the
corresponding section in :ref:`AssetOrTimeSchedule <asset-timetable-section>`.
+For more detailed information on asset-aware timetables, refer to
:ref:`AssetOrTimeSchedule <asset-timetable-section>`.
diff --git a/airflow-core/docs/authoring-and-scheduling/timetable.rst
b/airflow-core/docs/authoring-and-scheduling/timetable.rst
index 2deefe9429c..d516b731d5f 100644
--- a/airflow-core/docs/authoring-and-scheduling/timetable.rst
+++ b/airflow-core/docs/authoring-and-scheduling/timetable.rst
@@ -275,7 +275,10 @@ AssetOrTimeSchedule
Combining conditional asset expressions with time-based schedules enhances
scheduling flexibility.
-The ``AssetOrTimeSchedule`` is a specialized timetable that allows for the
scheduling of Dags based on both time-based schedules and asset events. It also
facilitates the creation of both scheduled runs, as per traditional timetables,
and asset-triggered runs, which operate independently.
+Asset-aware timetables let you combine a time-based schedule with an asset
expression:
+
+* ``AssetOrTimeSchedule`` schedules Dag runs both on the timetable and
whenever the assets update. It creates traditional scheduled runs and
asset-triggered runs independently.
+* ``AssetAndTimeSchedule`` keeps the Dag on a time-based timetable but only
creates a scheduled run after all referenced assets are ready. When the run is
created, the asset events are consumed so the next scheduled run waits for the
next set of updates. No asset-triggered runs are created.
This feature is particularly useful in scenarios where a Dag needs to run on
asset updates and also at periodic intervals. It ensures that the workflow
remains responsive to data changes and consistently runs regular checks or
updates.
@@ -283,8 +286,7 @@ Here's an example of a Dag using ``AssetOrTimeSchedule``:
.. code-block:: python
- from airflow.timetables.assets import AssetOrTimeSchedule
- from airflow.timetables.trigger import CronTriggerTimetable
+ from airflow.sdk import AssetOrTimeSchedule, CronTriggerTimetable
@dag(
@@ -296,6 +298,23 @@ Here's an example of a Dag using ``AssetOrTimeSchedule``:
def example_dag():
pass
+Here's an example of a Dag using ``AssetAndTimeSchedule`` to require both the
time-based schedule and fresh assets before a run is created:
+
+.. code-block:: python
+
+ from airflow.sdk import AssetAndTimeSchedule, CronTriggerTimetable
+
+
+ @dag(
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("0 1 * * 3", timezone="UTC"),
+ assets=(dag1_asset & dag2_asset),
+ ),
+ )
+ def example_gated_dag():
+ # Dag tasks go here
+ pass
+
Timetables comparisons
----------------------
diff --git a/airflow-core/docs/migrations-ref.rst
b/airflow-core/docs/migrations-ref.rst
index 4f5bf0f33ad..afd3ed01e8c 100644
--- a/airflow-core/docs/migrations-ref.rst
+++ b/airflow-core/docs/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
|
+=========================+==================+===================+==============================================================+
-| ``e5a91c7f42b3`` (head) | ``ca8499dc1004`` | ``3.4.0`` | Add
language column to dag_code. |
+| ``90e4d18ccadf`` (head) | ``e5a91c7f42b3`` | ``3.4.0`` | Add
timetable_asset_gated to DagModel. |
++-------------------------+------------------+-------------------+--------------------------------------------------------------+
+| ``e5a91c7f42b3`` | ``ca8499dc1004`` | ``3.4.0`` | Add
language column to dag_code. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
| ``ca8499dc1004`` | ``a61f0c9d2b47`` | ``3.4.0`` | Add
source_reference to import_error. |
+-------------------------+------------------+-------------------+--------------------------------------------------------------+
diff --git a/airflow-core/newsfragments/58543.feature.rst
b/airflow-core/newsfragments/58543.feature.rst
new file mode 100644
index 00000000000..da80079a161
--- /dev/null
+++ b/airflow-core/newsfragments/58543.feature.rst
@@ -0,0 +1 @@
+Add ``AssetAndTimeSchedule`` timetable that schedules time-based runs gated on
asset conditions.
diff --git a/airflow-core/src/airflow/dag_processing/collection.py
b/airflow-core/src/airflow/dag_processing/collection.py
index 33b1b443931..46099badce7 100644
--- a/airflow-core/src/airflow/dag_processing/collection.py
+++ b/airflow-core/src/airflow/dag_processing/collection.py
@@ -783,6 +783,7 @@ class DagModelOperation(NamedTuple):
dm.timetable_description = dag.timetable.description
dm.timetable_partitioned = dag.timetable.partitioned
dm.timetable_periodic = dag.timetable.periodic
+ dm.timetable_asset_gated = dag.timetable.asset_gated
dm.partition_mapper_info = dag.timetable.partition_mapper_info
dm.fail_fast = dag.fail_fast if dag.fail_fast is not None else
False
diff --git a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
index 741a7b4be6f..69819626e8c 100644
--- a/airflow-core/src/airflow/jobs/scheduler_job_runner.py
+++ b/airflow-core/src/airflow/jobs/scheduler_job_runner.py
@@ -26,7 +26,7 @@ import signal
import sys
import time
from collections import Counter, defaultdict, deque
-from collections.abc import Callable, Collection, Iterable, Iterator
+from collections.abc import Callable, Collection, Iterable, Iterator, Sequence
from contextlib import ExitStack, suppress
from datetime import datetime, timedelta
from functools import lru_cache, partial
@@ -117,7 +117,6 @@ from airflow.serialization.definitions.assets import
SerializedAssetUniqueKey
from airflow.serialization.definitions.notset import NOTSET
from airflow.ti_deps.dependencies_states import ACTIVE_STATES, EXECUTION_STATES
from airflow.timetables.base import Timetable, compute_rollup_fingerprint
-from airflow.timetables.simple import AssetTriggeredTimetable
from airflow.triggers.base import TriggerEvent
from airflow.utils.event_scheduler import EventScheduler
from airflow.utils.helpers import prune_dict
@@ -2711,6 +2710,19 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
)
continue
+ consumed_asset_records: Sequence[AssetDagRunQueue] = ()
+ gated_asset_events: list[AssetEvent] = []
+ if dag_model.timetable_asset_gated:
+ # The asset condition was evaluated without locks when this
Dag was
+ # selected; re-evaluate under ADRQ row locks so concurrent
schedulers
+ # cannot consume the same events twice. If the condition is
not (or no
+ # longer) satisfied, skip without touching the pending
schedule slot so
+ # a later loop retries it.
+ gate = self._collect_gated_asset_events(dag=serdag,
session=session)
+ if gate is None:
+ continue
+ consumed_asset_records, gated_asset_events = gate
+
try:
next_info =
serdag.timetable.next_run_info_from_dag_model(dag_model=dag_model)
if TYPE_CHECKING:
@@ -2744,6 +2756,11 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
session=session,
active_non_backfill_runs=active_runs_of_dags[dag_model.dag_id],
)
+ if consumed_asset_records:
+
created_run.consumed_asset_events.extend(gated_asset_events)
+ self._delete_consumed_asset_records(
+ records=consumed_asset_records,
dag_id=dag_model.dag_id, session=session
+ )
# Exceptions like ValueError, ParamValidationError, etc. are
raised by
# DagModel.create_dagrun() when dag is misconfigured. The
scheduler should not
@@ -2759,6 +2776,116 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
# 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,
load_assets=True, 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, load_assets: bool, session: Session
+ ) -> Sequence[AssetDagRunQueue]:
+ """Lock and return the Dag's queued asset (ADRQ) rows, skipping rows
another scheduler holds."""
+ query = select(AssetDagRunQueue).where(AssetDagRunQueue.target_dag_id
== dag_id)
+ if load_assets:
+ query = query.options(joinedload(AssetDagRunQueue.asset))
+ return session.scalars(
+ with_row_locks(
+ query,
+ of=AssetDagRunQueue,
+ skip_locked=True,
+ key_share=False,
+ session=session,
+ )
+ ).all()
+
+ def _select_consumed_asset_events(
+ self,
+ *,
+ dag: SerializedDAG,
+ records: Sequence[AssetDagRunQueue],
+ session: Session,
+ ) -> list[AssetEvent]:
+ """Select unconsumed events referenced by a Dag's locked ADRQ rows."""
+ referenced_event_ids = {record.asset_event_id for record in records}
+ event_predicate: ColumnElement[bool] =
AssetEvent.id.in_(referenced_event_ids)
+ if dag.catchup:
+ event_predicate = or_(
+ event_predicate,
+ AssetEvent.asset_id.in_(
+ select(DagScheduleAssetReference.asset_id).where(
+ DagScheduleAssetReference.dag_id == dag.dag_id
+ )
+ ),
+ AssetEvent.source_aliases.any(
+
AssetAliasModel.scheduled_dags.any(DagScheduleAssetAliasReference.dag_id ==
dag.dag_id)
+ ),
+ )
+ return list(
+ session.scalars(
+ select(AssetEvent)
+ .where(
+ event_predicate,
+ ~(
+ select(association_table.c.event_id)
+ .join(DagRun, DagRun.id ==
association_table.c.dag_run_id)
+ .where(
+ DagRun.dag_id == dag.dag_id,
+ association_table.c.event_id == AssetEvent.id,
+ )
+ .exists()
+ ),
+ )
+ .order_by(AssetEvent.timestamp.asc(), AssetEvent.id.asc())
+ )
+ )
+
+ def _delete_consumed_asset_records(
+ self, *, records: Sequence[AssetDagRunQueue], dag_id: str, session:
Session
+ ) -> None:
+ # Delete only consumed ADRQ rows to avoid dropping newly queued events
+ # (e.g. DagRun triggered by asset A while a new event for asset B
arrives).
+ result = cast(
+ "CursorResult",
+ session.execute(
+ delete(AssetDagRunQueue).where(
+ tuple_(
+ AssetDagRunQueue.target_dag_id,
+ AssetDagRunQueue.asset_event_id,
+ ).in_((record.target_dag_id, record.asset_event_id) for
record in records)
+ )
+ ),
+ )
+ self.log.info(
+ "Deleted %d ADRQ rows for '%s'",
+ result.rowcount,
+ dag_id,
+ )
+
def _create_dag_runs_asset_triggered(
self,
*,
@@ -2772,22 +2899,16 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
self.log.error("Dag '%s' not found in serialized_dag table",
dag_model.dag_id)
continue
- if not isinstance(dag.timetable, AssetTriggeredTimetable):
+ if not dag.timetable.asset_triggered:
self.log.error(
- "Dag '%s' was asset-scheduled, but didn't have an
AssetTriggeredTimetable!",
+ "Dag '%s' was routed to asset-triggered run creation, but
its timetable is not asset-triggered",
dag_model.dag_id,
)
continue
- queued_adrqs = session.scalars(
- with_row_locks(
-
select(AssetDagRunQueue).where(AssetDagRunQueue.target_dag_id == dag.dag_id),
- of=AssetDagRunQueue,
- skip_locked=True,
- key_share=False,
- session=session,
- )
- ).all()
+ queued_adrqs = self._lock_queued_asset_records(
+ dag_id=dag.dag_id, load_assets=False, session=session
+ )
# If another scheduler already locked these ADRQ rows, SKIP LOCKED
makes this scheduler skip them.
if not queued_adrqs:
self.log.debug(
@@ -2796,43 +2917,10 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
)
continue
- referenced_event_ids = {adrq.asset_event_id for adrq in
queued_adrqs}
- event_predicate: ColumnElement[bool] =
AssetEvent.id.in_(referenced_event_ids)
- if dag.catchup:
- # With catchup on, also consume events recorded before the Dag
started
- # scheduling on its assets/aliases, not just those with a
queue row. (With catchup
- # off only queued events are consumed.) The not-consumed
filter below dedupes
- # across runs, so no event window is needed.
- event_predicate = or_(
- event_predicate,
- AssetEvent.asset_id.in_(
- select(DagScheduleAssetReference.asset_id).where(
- DagScheduleAssetReference.dag_id == dag.dag_id
- )
- ),
- AssetEvent.source_aliases.any(
- AssetAliasModel.scheduled_dags.any(
- DagScheduleAssetAliasReference.dag_id == dag.dag_id
- )
- ),
- )
- asset_events = list(
- session.scalars(
- select(AssetEvent)
- .where(
- event_predicate,
- ~(
- select(association_table.c.event_id)
- .join(DagRun, DagRun.id ==
association_table.c.dag_run_id)
- .where(
- DagRun.dag_id == dag.dag_id,
- association_table.c.event_id == AssetEvent.id,
- )
- .exists()
- ),
- )
- .order_by(AssetEvent.timestamp.asc(), AssetEvent.id.asc())
- )
+ asset_events = self._select_consumed_asset_events(
+ dag=dag,
+ records=queued_adrqs,
+ session=session,
)
if asset_events:
triggered_date = timezone.coerce_datetime(max(event.timestamp
for event in asset_events))
@@ -2875,21 +2963,10 @@ class SchedulerJobRunner(BaseJobRunner, LoggingMixin):
)
# Always delete ADRQ rows for this batch to prevent stale entries
accumulating,
# including when all events were already consumed by a concurrent
DagRun.
- result = cast(
- "CursorResult",
- session.execute(
- delete(AssetDagRunQueue).where(
- tuple_(
- AssetDagRunQueue.target_dag_id,
- AssetDagRunQueue.asset_event_id,
- ).in_((adrq.target_dag_id, adrq.asset_event_id) for
adrq in queued_adrqs)
- )
- ),
- )
- self.log.info(
- "Deleted %d ADRQ rows for '%s'",
- result.rowcount,
- dag.dag_id,
+ self._delete_consumed_asset_records(
+ records=queued_adrqs,
+ dag_id=dag.dag_id,
+ session=session,
)
def _lock_backfills(self, dag_runs: Collection[DagRun], session: Session)
-> dict[int, Backfill]:
diff --git
a/airflow-core/src/airflow/migrations/versions/0141_3_4_0_add_timetable_asset_gated_to_dagmodel.py
b/airflow-core/src/airflow/migrations/versions/0141_3_4_0_add_timetable_asset_gated_to_dagmodel.py
new file mode 100644
index 00000000000..bd2be5e3a95
--- /dev/null
+++
b/airflow-core/src/airflow/migrations/versions/0141_3_4_0_add_timetable_asset_gated_to_dagmodel.py
@@ -0,0 +1,55 @@
+#
+# 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 timetable_asset_gated to DagModel.
+
+Revision ID: 90e4d18ccadf
+Revises: e5a91c7f42b3
+Create Date: 2026-06-21 06:42:26.369414
+
+"""
+
+from __future__ import annotations
+
+import sqlalchemy as sa
+from alembic import op
+
+from airflow.migrations.utils import disable_sqlite_fkeys
+
+revision = "90e4d18ccadf"
+down_revision = "e5a91c7f42b3"
+branch_labels = None
+depends_on = None
+airflow_version = "3.4.0"
+
+
+def upgrade():
+ """Add timetable_asset_gated column to DagModel."""
+ with disable_sqlite_fkeys(op):
+ with op.batch_alter_table("dag", schema=None) as batch_op:
+ batch_op.add_column(
+ sa.Column("timetable_asset_gated", sa.Boolean(),
server_default="0", nullable=False)
+ )
+
+
+def downgrade():
+ """Remove timetable_asset_gated column from DagModel."""
+ with disable_sqlite_fkeys(op):
+ with op.batch_alter_table("dag", schema=None) as batch_op:
+ batch_op.drop_column("timetable_asset_gated")
diff --git a/airflow-core/src/airflow/models/dag.py
b/airflow-core/src/airflow/models/dag.py
index 634910a858b..30107a0bee9 100644
--- a/airflow-core/src/airflow/models/dag.py
+++ b/airflow-core/src/airflow/models/dag.py
@@ -37,6 +37,7 @@ from sqlalchemy import (
Integer,
String,
Text,
+ and_,
case,
func,
inspect as sa_inspect,
@@ -71,7 +72,7 @@ from airflow.serialization.encoders import DAT,
encode_deadline_alert
from airflow.serialization.enums import Encoding
from airflow.timetables.base import DataInterval, PartitionMapperInfo,
Timetable
from airflow.timetables.interval import CronDataIntervalTimetable,
DeltaDataIntervalTimetable
-from airflow.timetables.simple import AssetTriggeredTimetable, NullTimetable,
OnceTimetable
+from airflow.timetables.simple import NullTimetable, OnceTimetable
from airflow.utils.session import NEW_SESSION, provide_session
from airflow.utils.sqlalchemy import UtcDateTime, with_row_locks
from airflow.utils.state import DagRunState, DagSchedulingState
@@ -131,8 +132,9 @@ def infer_automated_data_interval(timetable: Timetable,
logical_date: datetime)
:meta private:
"""
- timetable_type = type(timetable)
- if issubclass(timetable_type, (NullTimetable, OnceTimetable,
AssetTriggeredTimetable)):
+ if timetable.asset_triggered or issubclass(
+ timetable_type := type(timetable), (NullTimetable, OnceTimetable)
+ ):
return DataInterval.exact(timezone.coerce_datetime(logical_date))
start = timezone.coerce_datetime(logical_date)
if issubclass(timetable_type, CronDataIntervalTimetable):
@@ -362,6 +364,8 @@ class DagModel(Base):
timetable_partitioned: Mapped[bool] = mapped_column(Boolean,
nullable=False, server_default="0")
# Whether the timetable is periodic (supports backfilling).
timetable_periodic: Mapped[bool] = mapped_column(Boolean, nullable=False,
server_default="0")
+ # Whether the timetable's scheduled runs are gated on an asset condition.
+ timetable_asset_gated: Mapped[bool] = mapped_column(Boolean,
nullable=False, server_default="0")
# Cached partition mapper metadata for partitioned timetables, populated
# during Dag serialization so the UI can resolve mapper attributes without
# deserializing the timetable. See ``PartitionMapperInfo`` for the
per-asset
@@ -690,9 +694,10 @@ class DagModel(Base):
you should ensure that any scheduling decisions are made in a single
transaction -- as soon as the
transaction is committed it will be unlocked.
- For asset-triggered scheduling, Dags that have ``AssetDagRunQueue``
rows but no matching
- ``SerializedDagModel`` row are omitted from ``triggered_date_by_dag``
until serialization exists;
- ADRQs are **not** deleted here so the scheduler can re-evaluate on a
later run.
+ For asset-triggered and asset-gated scheduling, Dags that have
``AssetDagRunQueue`` rows
+ but no matching ``SerializedDagModel`` row are omitted from the
asset-aware scheduling
+ buckets until serialization exists; ADRQs are **not** deleted here so
the scheduler can
+ re-evaluate on a later run.
:meta private:
"""
@@ -731,7 +736,7 @@ class DagModel(Base):
if adrq_by_dag:
log.info(
- "Asset-triggered Dags with queued events: %s",
+ "Asset-aware Dags with queued events: %s",
{dag_id: len(adrqs) for dag_id, adrqs in adrq_by_dag.items()},
)
@@ -750,12 +755,19 @@ class DagModel(Base):
for dag_id in missing_from_serialized:
del adrq_by_dag[dag_id]
del dag_statuses[dag_id]
+ asset_gated_ready_dag_ids: set[str] = set()
for ser_dag in ser_dags:
dag_id = ser_dag.dag_id
statuses = dag_statuses[dag_id]
- ready = dag_ready(dag_id,
cond=ser_dag.dag.timetable.asset_condition, statuses=statuses)
+ timetable = ser_dag.dag.timetable
+ ready = dag_ready(dag_id, cond=timetable.asset_condition,
statuses=statuses)
if not ready:
log.debug("Asset condition not met for dag '%s'", dag_id)
+ if timetable.asset_gated and ready:
+ asset_gated_ready_dag_ids.add(dag_id)
+ if not (timetable.asset_triggered and ready):
+ # Only satisfied asset-triggered Dags stay in the
asset-triggered bucket
+ # (adrq_by_dag feeds triggered_date_by_dag below).
del adrq_by_dag[dag_id]
del dag_statuses[dag_id]
del dag_statuses
@@ -789,6 +801,7 @@ class DagModel(Base):
k: v for k, v in triggered_date_by_dag.items() if k not in
exclusion_list
}
+ time_due = cls.next_dagrun_create_after <= func.now()
# We limit so that _one_ scheduler doesn't try to do all the creation
of dag runs
query = (
select(cls)
@@ -799,8 +812,9 @@ class DagModel(Base):
cls.has_import_errors == expression.false(),
cls.exceeds_max_non_backfill == expression.false(),
or_(
- cls.next_dagrun_create_after <= func.now(),
cls.dag_id.in_(asset_triggered_dag_ids),
+ and_(cls.dag_id.in_(asset_gated_ready_dag_ids), time_due),
+ and_(cls.timetable_asset_gated == expression.false(),
time_due),
),
)
.order_by(cls.next_dagrun_create_after)
diff --git a/airflow-core/src/airflow/serialization/encoders.py
b/airflow-core/src/airflow/serialization/encoders.py
index 5f64ea0c76c..e4b417de213 100644
--- a/airflow-core/src/airflow/serialization/encoders.py
+++ b/airflow-core/src/airflow/serialization/encoders.py
@@ -34,6 +34,7 @@ from airflow.sdk import (
Asset,
AssetAlias,
AssetAll,
+ AssetAndTimeSchedule,
AssetAny,
AssetOrTimeSchedule,
ChainMapper,
@@ -323,6 +324,7 @@ class _Serializer:
"""Serialization logic."""
BUILTIN_TIMETABLES: dict[type, str] = {
+ AssetAndTimeSchedule: "airflow.timetables.assets.AssetAndTimeSchedule",
AssetOrTimeSchedule: "airflow.timetables.assets.AssetOrTimeSchedule",
AssetTriggeredTimetable:
"airflow.timetables.simple.AssetTriggeredTimetable",
ContinuousTimetable: "airflow.timetables.simple.ContinuousTimetable",
@@ -425,6 +427,13 @@ class _Serializer:
"run_immediately":
encode_run_immediately(representitive.run_immediately),
}
+ @serialize_timetable.register
+ def _(self, timetable: AssetAndTimeSchedule) -> dict[str, Any]:
+ return {
+ "asset_condition": encode_asset_like(timetable.asset_condition),
+ "timetable": encode_timetable(timetable.timetable),
+ }
+
@serialize_timetable.register
def _(self, timetable: AssetOrTimeSchedule) -> dict[str, Any]:
return {
diff --git a/airflow-core/src/airflow/timetables/assets.py
b/airflow-core/src/airflow/timetables/assets.py
index 37c6fb4825b..7359ba23d0e 100644
--- a/airflow-core/src/airflow/timetables/assets.py
+++ b/airflow-core/src/airflow/timetables/assets.py
@@ -18,18 +18,25 @@
from __future__ import annotations
import typing
+from collections.abc import Collection
from airflow.exceptions import AirflowTimetableInvalid
-from airflow.serialization.definitions.assets import SerializedAsset,
SerializedAssetBase
+from airflow.serialization.definitions.assets import SerializedAsset,
SerializedAssetAll, SerializedAssetBase
+from airflow.timetables.base import Timetable
from airflow.timetables.simple import AssetTriggeredTimetable
from airflow.utils.types import DagRunType
if typing.TYPE_CHECKING:
- from collections.abc import Collection
-
import pendulum
- from airflow.timetables.base import DagRunInfo, DataInterval,
TimeRestriction, Timetable
+ from airflow.timetables.base import DagRunInfo, DataInterval,
TimeRestriction
+
+
+def _validate_asset_time_schedule(*, timetable: Timetable, asset_condition:
SerializedAssetBase) -> None:
+ if timetable.asset_triggered or timetable.asset_gated:
+ raise AirflowTimetableInvalid("Cannot nest asset-aware timetables")
+ if not isinstance(asset_condition, SerializedAssetBase):
+ raise AirflowTimetableInvalid("All elements in 'assets' must be
assets")
class AssetOrTimeSchedule(AssetTriggeredTimetable):
@@ -58,10 +65,10 @@ class AssetOrTimeSchedule(AssetTriggeredTimetable):
)
def validate(self) -> None:
- if isinstance(self.timetable, AssetTriggeredTimetable):
- raise AirflowTimetableInvalid("cannot nest asset timetables")
- if not isinstance(self.asset_condition, SerializedAssetBase):
- raise AirflowTimetableInvalid("all elements in 'assets' must be
assets")
+ _validate_asset_time_schedule(
+ timetable=self.timetable,
+ asset_condition=self.asset_condition,
+ )
def serialize(self) -> dict[str, typing.Any]:
from airflow.serialization.encoders import encode_asset_like,
encode_timetable
@@ -90,3 +97,92 @@ class AssetOrTimeSchedule(AssetTriggeredTimetable):
if run_type != DagRunType.ASSET_TRIGGERED:
return self.timetable.generate_run_id(run_type=run_type, **kwargs)
return super().generate_run_id(run_type=run_type, **kwargs)
+
+
+class AssetAndTimeSchedule(Timetable):
+ """
+ Time-based schedule that waits for required assets before creating a run.
+
+ This timetable composes a time-based timetable with an asset condition. It
+ schedules runs according to the provided ``timetable`` (e.g. cron), but a
run
+ is only created when all required assets are present. Unlike
+ :class:`AssetOrTimeSchedule`, this does not create asset-triggered runs.
+ """
+
+ asset_gated = True
+
+ def __init__(
+ self,
+ *,
+ timetable: Timetable,
+ assets: Collection[SerializedAsset] | SerializedAssetBase,
+ ) -> None:
+ from airflow.serialization.encoders import ensure_serialized_asset
+
+ self.timetable = timetable
+
+ if isinstance(assets, SerializedAssetBase):
+ self.asset_condition = assets
+ elif isinstance(assets, Collection):
+ self.asset_condition =
SerializedAssetAll([ensure_serialized_asset(a) for a in assets])
+ else:
+ self.asset_condition = ensure_serialized_asset(assets)
+
+ @classmethod
+ def deserialize(cls, data: dict[str, typing.Any]) -> Timetable:
+ from airflow.serialization.decoders import decode_asset_like,
decode_timetable
+
+ return cls(
+ assets=decode_asset_like(data["asset_condition"]),
+ timetable=decode_timetable(data["timetable"]),
+ )
+
+ def serialize(self) -> dict[str, typing.Any]:
+ from airflow.serialization.encoders import encode_asset_like,
encode_timetable
+
+ return {
+ "asset_condition": encode_asset_like(self.asset_condition),
+ "timetable": encode_timetable(self.timetable),
+ }
+
+ def validate(self) -> None:
+ _validate_asset_time_schedule(
+ timetable=self.timetable,
+ asset_condition=self.asset_condition,
+ )
+
+ @property
+ def description(self) -> str: # type: ignore[override]
+ return f"Triggered by assets and {self.timetable.description}"
+
+ @property
+ def summary(self) -> str:
+ return f"Asset and {self.timetable.summary}"
+
+ @property
+ def periodic(self) -> bool: # type: ignore[override]
+ return self.timetable.periodic
+
+ @property
+ def can_be_scheduled(self) -> bool: # type: ignore[override]
+ return self.timetable.can_be_scheduled
+
+ @property
+ def active_runs_limit(self) -> int | None: # type: ignore[override]
+ return self.timetable.active_runs_limit
+
+ def infer_manual_data_interval(self, *, run_after: pendulum.DateTime) ->
DataInterval:
+ return self.timetable.infer_manual_data_interval(run_after=run_after)
+
+ def next_dagrun_info(
+ self, *, last_automated_data_interval: DataInterval | None,
restriction: TimeRestriction
+ ) -> DagRunInfo | None:
+ return self.timetable.next_dagrun_info(
+ last_automated_data_interval=last_automated_data_interval,
+ restriction=restriction,
+ )
+
+ def generate_run_id(self, *, run_type: DagRunType, **kwargs: typing.Any)
-> str:
+ # All run IDs are delegated to the wrapped timetable; this class
+ # intentionally does not create ASSET_TRIGGERED runs.
+ return self.timetable.generate_run_id(run_type=run_type, **kwargs)
diff --git a/airflow-core/src/airflow/timetables/base.py
b/airflow-core/src/airflow/timetables/base.py
index 365f980b932..49ebc50e288 100644
--- a/airflow-core/src/airflow/timetables/base.py
+++ b/airflow-core/src/airflow/timetables/base.py
@@ -237,6 +237,22 @@ class Timetable(Protocol):
instead of the traditional logic based on logical dates and data intervals.
"""
+ asset_triggered: bool = False
+ """Whether this timetable creates runs triggered by asset events.
+
+ This is *True* for timetables that materialize an asset-triggered DagRun as
+ soon as their asset condition is satisfied, independently of a time
+ schedule.
+ """
+
+ asset_gated: bool = False
+ """Whether this timetable's scheduled runs are gated on an asset condition.
+
+ This is *True* for timetables whose time-based runs are only created once
+ their asset condition is satisfied. Manual and backfill runs are
unaffected.
+ A timetable that enables this must define a non-null ``asset_condition``.
+ """
+
partitioned_at_runtime: bool = False
"""Whether this timetable defers partition selection to task runtime.
diff --git a/airflow-core/src/airflow/timetables/simple.py
b/airflow-core/src/airflow/timetables/simple.py
index 87741a8eb0d..a1c64a213b0 100644
--- a/airflow-core/src/airflow/timetables/simple.py
+++ b/airflow-core/src/airflow/timetables/simple.py
@@ -217,6 +217,7 @@ class AssetTriggeredTimetable(_TrivialTimetable):
"""
description: str = "Triggered by assets"
+ asset_triggered = True
def __init__(self, assets: Collection[SerializedAsset] |
SerializedAssetBase) -> None:
super().__init__()
diff --git a/airflow-core/src/airflow/utils/db.py
b/airflow-core/src/airflow/utils/db.py
index 641f87a1f87..f160cb0597d 100644
--- a/airflow-core/src/airflow/utils/db.py
+++ b/airflow-core/src/airflow/utils/db.py
@@ -117,7 +117,7 @@ _REVISION_HEADS_MAP: dict[str, str] = {
"3.1.8": "509b94a1042d",
"3.2.0": "1d6611b6ab7c",
"3.3.0": "d2f4e1b3c5a7",
- "3.4.0": "e5a91c7f42b3",
+ "3.4.0": "90e4d18ccadf",
}
# Prefix used to identify tables holding data moved during migration.
diff --git a/airflow-core/tests/unit/dag_processing/test_collection.py
b/airflow-core/tests/unit/dag_processing/test_collection.py
index cc331423e3f..da9342958d6 100644
--- a/airflow-core/tests/unit/dag_processing/test_collection.py
+++ b/airflow-core/tests/unit/dag_processing/test_collection.py
@@ -73,7 +73,14 @@ from airflow.partition_mappers.window import DayWindow
from airflow.plugins_manager import AirflowPlugin
from airflow.providers.standard.operators.empty import EmptyOperator
from airflow.providers.standard.triggers.file import FileDeleteTrigger
-from airflow.sdk import DAG, Asset, AssetAlias, AssetAll, AssetWatcher
+from airflow.sdk import (
+ DAG,
+ Asset,
+ AssetAlias,
+ AssetAll,
+ AssetAndTimeSchedule,
+ AssetWatcher,
+)
from airflow.sdk.definitions.deadline import AsyncCallback,
BaseDeadlineReference, DeadlineAlert
from airflow.sdk.definitions.timetables.assets import AssetOrTimeSchedule,
PartitionedAssetTimetable
from airflow.serialization.definitions.assets import SerializedAsset
@@ -913,6 +920,31 @@ class TestUpdateDagParsingResults:
dag_model: DagModel = session.get(DagModel, (dag.dag_id,))
assert dag_model.last_parse_duration == parse_duration
+ def test_timetable_asset_gated_written_to_db_on_sync(self,
testing_dag_bundle, session):
+ asset = Asset("test")
+ gated_dag = DAG(
+ dag_id="asset_gated",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("@daily", timezone="UTC"),
+ assets=asset,
+ ),
+ catchup=False,
+ )
+ regular_dag = DAG(dag_id="regular", schedule=None)
+
+ update_dag_parsing_results_in_db(
+ "testing",
+ None,
+ [LazyDeserializedDAG.from_dag(gated_dag),
LazyDeserializedDAG.from_dag(regular_dag)],
+ {},
+ None,
+ set(),
+ session,
+ )
+
+ assert session.get(DagModel, gated_dag.dag_id).timetable_asset_gated
is True
+ assert session.get(DagModel, regular_dag.dag_id).timetable_asset_gated
is False
+
@patch.object(ParseImportError, "full_file_path")
@patch.object(SerializedDagModel, "write_dag")
@pytest.mark.usefixtures("clean_db")
diff --git a/airflow-core/tests/unit/jobs/test_scheduler_job.py
b/airflow-core/tests/unit/jobs/test_scheduler_job.py
index af339fa0e92..b7244388e82 100644
--- a/airflow-core/tests/unit/jobs/test_scheduler_job.py
+++ b/airflow-core/tests/unit/jobs/test_scheduler_job.py
@@ -29,7 +29,8 @@ from concurrent.futures import ThreadPoolExecutor,
as_completed
from contextlib import ExitStack, contextmanager
from datetime import timedelta
from pathlib import Path
-from typing import TYPE_CHECKING
+from types import SimpleNamespace
+from typing import TYPE_CHECKING, cast
from unittest import mock
from unittest.mock import MagicMock, patch
from uuid import UUID, uuid4
@@ -125,8 +126,10 @@ from airflow.sdk import (
DAG,
Asset,
AssetAlias,
+ AssetAndTimeSchedule,
AssetWatcher,
CronPartitionTimetable,
+ CronTriggerTimetable,
FixedKeyMapper,
HourWindow,
IdentityMapper,
@@ -142,13 +145,13 @@ from airflow.sdk.definitions.timetables.assets import
PartitionedAssetTimetable
from airflow.serialization.definitions.dag import SerializedDAG
from airflow.serialization.encoders import ensure_serialized_asset
from airflow.serialization.serialized_objects import LazyDeserializedDAG
-from airflow.timetables.base import DagRunInfo, DataInterval,
compute_rollup_fingerprint
+from airflow.timetables.base import DagRunInfo, DataInterval, Timetable,
compute_rollup_fingerprint
from airflow.timetables.simple import (
PartitionedAssetTimetable as CorePartitionedAssetTimetable,
PartitionedAtRuntime,
)
from airflow.utils.session import NEW_SESSION, create_session, provide_session
-from airflow.utils.sqlalchemy import with_row_locks
+from airflow.utils.sqlalchemy import CommitProhibitorGuard, with_row_locks
from airflow.utils.state import CallbackState, DagRunState,
DagSchedulingState, State, TaskInstanceState
from airflow.utils.types import DagRunTriggeredByType, DagRunType
@@ -6057,7 +6060,7 @@ class TestSchedulerJob:
self.job_runner = SchedulerJobRunner(job=scheduler_job,
executors=[self.null_exec])
with create_session() as session:
- self.job_runner._create_dagruns_for_dags(session, session)
+
self.job_runner._create_dagruns_for_dags(cast("CommitProhibitorGuard",
session), session)
def dict_from_obj(obj):
"""Get dict of column attrs from SqlAlchemy object."""
@@ -6094,6 +6097,23 @@ class TestSchedulerJob:
assert created_run.creating_job_id == scheduler_job.id
+ @mock.patch.object(SchedulerJobRunner, "_get_current_dag", autospec=True)
+ def test_create_dag_runs_asset_triggered_uses_behavior_flag(self,
mock_get_current_dag):
+ dag_model = SimpleNamespace(dag_id="custom-asset-triggered")
+ dag = SimpleNamespace(
+ dag_id=dag_model.dag_id,
+ timetable=SimpleNamespace(asset_triggered=True),
+ )
+ mock_get_current_dag.return_value = dag
+ session = MagicMock(spec=["get_bind", "scalars"])
+ session.get_bind.return_value.dialect.name = "postgresql"
+ session.scalars.return_value.all.return_value = []
+ runner = SchedulerJobRunner(job=Job(), executors=[self.null_exec])
+
+ runner._create_dag_runs_asset_triggered(dag_models=[dag_model],
session=session)
+
+ session.scalars.assert_called_once()
+
@pytest.mark.need_serialized_dag
@pytest.mark.parametrize(
("catchup", "expects_old_event"),
@@ -6466,7 +6486,7 @@ class TestSchedulerJob:
self.job_runner = SchedulerJobRunner(job=scheduler_job,
executors=[self.null_exec])
with create_session() as session:
- self.job_runner._create_dagruns_for_dags(session, session)
+
self.job_runner._create_dagruns_for_dags(cast("CommitProhibitorGuard",
session), session)
def dict_from_obj(obj):
"""Get dict of column attrs from SqlAlchemy object."""
@@ -11427,6 +11447,542 @@ def
test_schedule_dag_run_with_upstream_skip(dag_maker, session):
assert running_count == 2
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def test_create_dagruns_asset_and_time_waits_until_assets_ready(session:
Session, dag_maker):
+ asset = Asset(uri="test://asset-and-time-waits",
name="asset-and-time-waits")
+ logical_date = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-waits",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("* * * * *", timezone="UTC"),
+ assets=[asset],
+ ),
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun = logical_date
+ dag_model.next_dagrun_data_interval = (logical_date, logical_date)
+ dag_model.next_dagrun_create_after = logical_date
+ session.flush()
+
+ SchedulerJobRunner(job=Job(),
executors=[MockExecutor()])._create_dagruns_for_dags(
+ cast("CommitProhibitorGuard", session), session
+ )
+ session.flush()
+
+ assert session.scalar(select(DagRun).where(DagRun.dag_id ==
dag_model.dag_id)) is None
+ session.refresh(dag_model)
+ assert dag_model.next_dagrun == logical_date
+ assert dag_model.next_dagrun_create_after == logical_date
+
+
[email protected](
+ ("asset_triggered", "asset_gated"),
+ [
+ pytest.param(True, False, id="asset-triggered"),
+ pytest.param(False, True, id="asset-gated"),
+ ],
+)
[email protected](SerializedDagModel, "get_latest_serialized_dags",
autospec=True)
+def test_dags_needing_dagruns_routes_custom_timetable_by_behavior(
+ mock_get_latest_serialized_dags, asset_triggered, asset_gated, session:
Session, dag_maker
+):
+ asset = Asset(uri="test://custom-asset-scheduling",
name="custom-asset-scheduling")
+ with dag_maker(
+ dag_id=f"custom-asset-scheduling-{asset_triggered}-{asset_gated}",
+ schedule=[asset],
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun_create_after = timezone.utcnow() -
timedelta(minutes=1)
+ dag_model.timetable_asset_gated = asset_gated
+
+ asset_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset.uri))
+ asset_event = AssetEvent(asset_id=asset_id, timestamp=timezone.utcnow())
+ session.add(asset_event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=asset_event.id,
+ )
+ )
+ session.flush()
+
+ timetable = MagicMock(spec=Timetable)
+ timetable.asset_triggered = asset_triggered
+ timetable.asset_gated = asset_gated
+ timetable.asset_condition = ensure_serialized_asset(asset)
+ serialized_dag = SimpleNamespace(
+ dag_id=dag_model.dag_id,
+ dag=SimpleNamespace(timetable=timetable),
+ )
+ mock_get_latest_serialized_dags.return_value = [serialized_dag]
+
+ query, triggered_date_by_dag = DagModel.dags_needing_dagruns(session)
+
+ # Both behaviors keep the Dag selected for run creation; only
asset-triggered
+ # timetables land in the asset-triggered bucket (gated Dags take the normal
+ # scheduled path). The gated Dag is selected via its satisfied asset
condition:
+ # with timetable_asset_gated=True, being time-due alone would not select
it.
+ assert [model.dag_id for model in query.all()] == [dag_model.dag_id]
+ assert (dag_model.dag_id in triggered_date_by_dag) is asset_triggered
+
+
+@time_machine.travel("2026-03-29 18:30:00+00:00")
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def
test_create_dagruns_asset_and_time_late_arrival_uses_oldest_pending_logical_date(
+ session: Session, dag_maker
+):
+ asset = Asset(uri="test://asset-and-time-ready",
name="asset-and-time-ready")
+ logical_date = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ asset_created_at = pendulum.datetime(2026, 3, 29, 18, 30, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-ready",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("* * * * *", timezone="UTC"),
+ assets=[asset],
+ ),
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun = logical_date
+ dag_model.next_dagrun_data_interval = (logical_date, logical_date)
+ dag_model.next_dagrun_create_after = logical_date
+
+ asset_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset.uri))
+ asset_event = AssetEvent(asset_id=asset_id, timestamp=asset_created_at)
+ session.add(asset_event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=asset_event.id,
+ created_at=asset_created_at,
+ )
+ )
+ session.flush()
+
+ SchedulerJobRunner(job=Job(),
executors=[MockExecutor()])._create_dagruns_for_dags(
+ cast("CommitProhibitorGuard", session), session
+ )
+ session.flush()
+
+ dag_run = session.scalars(select(DagRun).where(DagRun.dag_id ==
dag_model.dag_id)).one()
+ assert dag_run.state == DagRunState.QUEUED
+ assert dag_run.run_type == DagRunType.SCHEDULED
+ assert dag_run.logical_date == logical_date
+ assert (
+
session.scalar(select(AssetDagRunQueue).where(AssetDagRunQueue.target_dag_id ==
dag_model.dag_id))
+ is None
+ )
+
+
+@time_machine.travel("2026-03-29 19:00:00+00:00")
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def
test_create_dagruns_asset_and_time_late_arrival_consumes_only_one_slot(session:
Session, dag_maker):
+ asset = Asset(uri="test://asset-and-time-one-slot",
name="asset-and-time-one-slot")
+ first_logical_date = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ second_logical_date = pendulum.datetime(2026, 3, 29, 18, tz="UTC")
+ asset_created_at = pendulum.datetime(2026, 3, 29, 19, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-one-slot",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("0 * * * *", timezone="UTC"),
+ assets=[asset],
+ ),
+ catchup=True,
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun = first_logical_date
+ dag_model.next_dagrun_data_interval = (first_logical_date,
first_logical_date)
+ dag_model.next_dagrun_create_after = first_logical_date
+
+ asset_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset.uri))
+ asset_event = AssetEvent(
+ asset_id=asset_id,
+ source_task_id="produce",
+ source_dag_id="producer",
+ source_run_id="producer_run",
+ source_map_index=-1,
+ timestamp=asset_created_at,
+ )
+ session.add(asset_event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=asset_event.id,
+ created_at=asset_created_at,
+ )
+ )
+ session.flush()
+
+ SchedulerJobRunner(job=Job(),
executors=[MockExecutor()])._create_dagruns_for_dags(
+ cast("CommitProhibitorGuard", session), session
+ )
+ session.flush()
+
+ dag_runs = session.scalars(select(DagRun).where(DagRun.dag_id ==
dag_model.dag_id)).all()
+ assert len(dag_runs) == 1
+ assert dag_runs[0].run_type == DagRunType.SCHEDULED
+ assert dag_runs[0].logical_date == first_logical_date
+ session.refresh(dag_model)
+ assert dag_model.next_dagrun == second_logical_date
+ assert dag_model.next_dagrun_create_after == second_logical_date
+
+
+@time_machine.travel("2026-03-29 17:30:00+00:00")
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def test_create_dagruns_asset_and_time_respects_max_active_runs(session:
Session, dag_maker):
+ """
+ Regression test: creating an asset-gated run must update
exceeds_max_non_backfill
+ so that a follow-up asset event does not bypass max_active_runs. The SQL
filter in
+ dags_needing_dagruns relies on DagModel.exceeds_max_non_backfill; if run
creation
+ skipped _set_exceeds_max_active_runs, the next loop would create a second
run
+ even though max_active_runs=1.
+ """
+ asset = Asset(uri="test://asset-and-time-max-active",
name="asset-and-time-max-active")
+ first_logical_date = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ first_event_at = pendulum.datetime(2026, 3, 29, 17, 0, 30, tz="UTC")
+ second_event_at = pendulum.datetime(2026, 3, 29, 17, 1, 30, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-max-active",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("* * * * *", timezone="UTC"),
+ assets=[asset],
+ ),
+ max_active_runs=1,
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun = first_logical_date
+ dag_model.next_dagrun_data_interval = (first_logical_date,
first_logical_date)
+ dag_model.next_dagrun_create_after = first_logical_date
+
+ asset_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset.uri))
+ first_event = AssetEvent(asset_id=asset_id, timestamp=first_event_at)
+ session.add(first_event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=first_event.id,
+ created_at=first_event_at,
+ )
+ )
+ session.flush()
+
+ job_runner = SchedulerJobRunner(job=Job(), executors=[MockExecutor()])
+ job_runner._create_dagruns_for_dags(cast("CommitProhibitorGuard",
session), session)
+ session.flush()
+
+ first_runs = session.scalars(select(DagRun).where(DagRun.dag_id ==
dag_model.dag_id)).all()
+ assert len(first_runs) == 1
+ session.refresh(dag_model)
+ # After creating the first run with max_active_runs=1, the DagModel must
be flagged
+ # as exceeding max active runs so dags_needing_dagruns excludes it in the
next loop.
+ assert dag_model.exceeds_max_non_backfill is True
+
+ # Simulate next producer event: new ADRQ written while the first run is
still queued.
+ second_event = AssetEvent(asset_id=asset_id, timestamp=second_event_at)
+ session.add(second_event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=second_event.id,
+ created_at=second_event_at,
+ )
+ )
+ session.flush()
+
+ job_runner._create_dagruns_for_dags(cast("CommitProhibitorGuard",
session), session)
+ session.flush()
+
+ runs_after_second_loop =
session.scalars(select(DagRun).where(DagRun.dag_id == dag_model.dag_id)).all()
+ # Second loop must NOT create a new run because max_active_runs=1 is
already saturated.
+ assert len(runs_after_second_loop) == 1
+
+
+@time_machine.travel("2026-03-29 17:30:00+00:00")
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def
test_create_dagruns_asset_and_time_populates_consumed_asset_events(session:
Session, dag_maker):
+ """
+ Regression test: asset-gated runs must carry the consumed AssetEvent rows
on
+ DagRun.consumed_asset_events so that triggering_asset_events templates,
+ inlet_events callbacks, and the UI asset provenance section work the same
way
+ as asset-triggered runs do. Like asset-triggered runs with catchup off,
events
+ that predate the Dag scheduling on its assets are backlog and must be
skipped.
+ """
+ asset = Asset(uri="test://asset-and-time-consumed",
name="asset-and-time-consumed")
+ logical_date = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ registered_at = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ backlog_event_at = pendulum.datetime(2026, 3, 29, 16, 30, tz="UTC")
+ event_at = pendulum.datetime(2026, 3, 29, 17, 0, 30, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-consumed",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("* * * * *", timezone="UTC"),
+ assets=[asset],
+ ),
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun = logical_date
+ dag_model.next_dagrun_data_interval = (logical_date, logical_date)
+ dag_model.next_dagrun_create_after = logical_date
+ session.scalars(
+
select(DagScheduleAssetReference).where(DagScheduleAssetReference.dag_id ==
dag_model.dag_id)
+ ).one().created_at = registered_at
+
+ asset_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset.uri))
+ backlog_asset_event = AssetEvent(
+ asset_id=asset_id,
+ source_task_id="produce",
+ source_dag_id="producer",
+ source_run_id="producer_backlog_run",
+ source_map_index=-1,
+ timestamp=backlog_event_at,
+ )
+ asset_event = AssetEvent(
+ asset_id=asset_id,
+ source_task_id="produce",
+ source_dag_id="producer",
+ source_run_id="producer_run",
+ source_map_index=-1,
+ timestamp=event_at,
+ )
+ session.add_all([backlog_asset_event, asset_event])
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=asset_event.id,
+ created_at=event_at,
+ )
+ )
+ session.flush()
+
+ SchedulerJobRunner(job=Job(),
executors=[MockExecutor()])._create_dagruns_for_dags(
+ cast("CommitProhibitorGuard", session), session
+ )
+ session.flush()
+
+ dag_run = session.scalars(select(DagRun).where(DagRun.dag_id ==
dag_model.dag_id)).one()
+ # The asset event that satisfied the gate must be linked to the run for
+ # provenance (triggering_asset_events template, callback context, UI); the
+ # event from before the Dag scheduled on the asset must not.
+ assert list(dag_run.consumed_asset_events) == [asset_event]
+
+
+@time_machine.travel("2026-03-29 19:00:00+00:00")
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def
test_create_dagruns_asset_and_time_does_not_reattribute_consumed_events(session:
Session, dag_maker):
+ """
+ Regression test: each queued event must appear on exactly one run's
+ consumed_asset_events, even when catchup includes unconsumed backlog
events.
+ """
+ asset = Asset(uri="test://asset-and-time-no-reattribute",
name="asset-and-time-no-reattribute")
+ first_slot = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ first_event_at = pendulum.datetime(2026, 3, 29, 17, 2, tz="UTC")
+ second_event_at = pendulum.datetime(2026, 3, 29, 18, 2, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-no-reattribute",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("0 * * * *", timezone="UTC"),
+ assets=[asset],
+ ),
+ catchup=True,
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun = first_slot
+ dag_model.next_dagrun_data_interval = (first_slot, first_slot)
+ dag_model.next_dagrun_create_after = first_slot
+
+ asset_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset.uri))
+ first_event = AssetEvent(
+ asset_id=asset_id,
+ source_task_id="produce",
+ source_dag_id="producer",
+ source_run_id="producer_run_1",
+ source_map_index=-1,
+ timestamp=first_event_at,
+ )
+ session.add(first_event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=first_event.id,
+ created_at=first_event_at,
+ )
+ )
+ session.flush()
+
+ job_runner = SchedulerJobRunner(job=Job(), executors=[MockExecutor()])
+ job_runner._create_dagruns_for_dags(cast("CommitProhibitorGuard",
session), session)
+ session.flush()
+
+ second_event = AssetEvent(
+ asset_id=asset_id,
+ source_task_id="produce",
+ source_dag_id="producer",
+ source_run_id="producer_run_2",
+ source_map_index=-1,
+ timestamp=second_event_at,
+ )
+ session.add(second_event)
+ session.flush()
+ session.add(
+ AssetDagRunQueue(
+ asset_id=asset_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=second_event.id,
+ created_at=second_event_at,
+ )
+ )
+ session.flush()
+
+ job_runner._create_dagruns_for_dags(cast("CommitProhibitorGuard",
session), session)
+ session.flush()
+
+ runs = session.scalars(
+ select(DagRun).where(DagRun.dag_id ==
dag_model.dag_id).order_by(DagRun.logical_date)
+ ).all()
+ assert len(runs) == 2
+ assert list(runs[0].consumed_asset_events) == [first_event]
+ assert list(runs[1].consumed_asset_events) == [second_event]
+
+
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def test_create_dagruns_asset_and_time_rechecks_locked_adrq_rows(session:
Session, dag_maker):
+ asset_1 = Asset(uri="test://asset-and-time-locked-1",
name="asset-and-time-locked-1")
+ asset_2 = Asset(uri="test://asset-and-time-locked-2",
name="asset-and-time-locked-2")
+ logical_date = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-locked",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("* * * * *", timezone="UTC"),
+ assets=asset_1 & asset_2,
+ ),
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ dag_model = dag_maker.dag_model
+ dag_model.next_dagrun = logical_date
+ dag_model.next_dagrun_data_interval = (logical_date, logical_date)
+ dag_model.next_dagrun_create_after = logical_date
+ asset_1_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset_1.uri))
+ asset_2_id = session.scalar(select(AssetModel.id).where(AssetModel.uri ==
asset_2.uri))
+ event_1 = AssetEvent(asset_id=asset_1_id, timestamp=timezone.utcnow())
+ event_2 = AssetEvent(asset_id=asset_2_id, timestamp=timezone.utcnow())
+ session.add_all([event_1, event_2])
+ session.flush()
+ session.add_all(
+ [
+ AssetDagRunQueue(
+ asset_id=asset_1_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=event_1.id,
+ created_at=timezone.utcnow(),
+ ),
+ AssetDagRunQueue(
+ asset_id=asset_2_id,
+ target_dag_id=dag_model.dag_id,
+ asset_event_id=event_2.id,
+ created_at=timezone.utcnow(),
+ ),
+ ]
+ )
+ session.flush()
+
+ job_runner = SchedulerJobRunner(job=Job(), executors=[MockExecutor()])
+
+ def _lock_only_selected_row(query, **_):
+ if query.column_descriptions and
query.column_descriptions[0].get("entity") is AssetDagRunQueue:
+ return query.where(AssetDagRunQueue.asset_id == asset_1_id)
+ return query
+
+ with patch("airflow.jobs.scheduler_job_runner.with_row_locks",
side_effect=_lock_only_selected_row):
+ job_runner._create_dagruns_for_dags(cast("CommitProhibitorGuard",
session), session)
+
+ assert session.scalar(select(DagRun).where(DagRun.dag_id ==
dag_model.dag_id)) is None
+
+ remaining_adrq_asset_ids = set(
+ session.scalars(
+
select(AssetDagRunQueue.asset_id).where(AssetDagRunQueue.target_dag_id ==
dag_model.dag_id)
+ )
+ )
+ assert remaining_adrq_asset_ids == {asset_1_id, asset_2_id}
+
+
+@time_machine.travel("2026-03-29 22:40:00+00:00")
[email protected]("disable_load_example")
[email protected]_serialized_dag
+def
test_dags_needing_dagruns_asset_and_time_missing_assets_do_not_starve_time_dags(
+ session: Session, dag_maker
+):
+ asset = Asset(uri="test://asset-and-time-starvation",
name="asset-and-time-starvation")
+ gated_logical_date = pendulum.datetime(2026, 3, 29, 17, tz="UTC")
+ time_logical_date = pendulum.datetime(2026, 3, 29, 18, tz="UTC")
+ with dag_maker(
+ dag_id="asset-and-time-starvation",
+ schedule=AssetAndTimeSchedule(
+ timetable=CronTriggerTimetable("* * * * *", timezone="UTC"),
+ assets=[asset],
+ ),
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ gated_dag_model = dag_maker.dag_model
+ gated_dag_model.next_dagrun = gated_logical_date
+ gated_dag_model.next_dagrun_data_interval = (gated_logical_date,
gated_logical_date)
+ gated_dag_model.next_dagrun_create_after = gated_logical_date
+
+ with dag_maker(
+ dag_id="pure-time-after-gated",
+ schedule=CronTriggerTimetable("* * * * *", timezone="UTC"),
+ session=session,
+ ):
+ EmptyOperator(task_id="dummy_task")
+ time_dag_model = dag_maker.dag_model
+ time_dag_model.next_dagrun = time_logical_date
+ time_dag_model.next_dagrun_data_interval = (time_logical_date,
time_logical_date)
+ time_dag_model.next_dagrun_create_after = time_logical_date
+ session.flush()
+
+ with mock.patch.object(DagModel, "NUM_DAGS_PER_DAGRUN_QUERY", 1):
+ query, _ = DagModel.dags_needing_dagruns(session)
+
+ # The gated Dag has no queued assets, so it must not occupy the (mocked to
1)
+ # query limit slot even though its slot sorts first; the time Dag gets it.
+ assert [dag_model.dag_id for dag_model in query.all()] ==
[time_dag_model.dag_id]
+
+
class TestSchedulerJobQueriesCount:
"""
These tests are designed to detect changes in the number of queries for
diff --git a/airflow-core/tests/unit/models/test_dag.py
b/airflow-core/tests/unit/models/test_dag.py
index 922dca841de..e821d3efe8d 100644
--- a/airflow-core/tests/unit/models/test_dag.py
+++ b/airflow-core/tests/unit/models/test_dag.py
@@ -57,6 +57,7 @@ from airflow.models.dag import (
clear_team_name_cache,
get_next_data_interval,
get_run_data_interval,
+ infer_automated_data_interval,
)
from airflow.models.dagbag import DBDagBag
from airflow.models.dagbundle import DagBundleModel
@@ -168,6 +169,17 @@ def test_dags_bundle(configure_testing_dag_bundle):
yield
+def test_infer_automated_data_interval_uses_asset_triggered_behavior():
+ class CustomAssetTriggeredTimetable(Timetable):
+ asset_triggered = True
+
+ logical_date = timezone.datetime(2026, 6, 21)
+
+ assert infer_automated_data_interval(CustomAssetTriggeredTimetable(),
logical_date) == DataInterval.exact(
+ logical_date
+ )
+
+
def _create_dagrun(
dag: DAG,
*,
diff --git a/airflow-core/tests/unit/timetables/test_assets_timetable.py
b/airflow-core/tests/unit/timetables/test_assets_timetable.py
index ccdc395ae41..5ab7667cd11 100644
--- a/airflow-core/tests/unit/timetables/test_assets_timetable.py
+++ b/airflow-core/tests/unit/timetables/test_assets_timetable.py
@@ -28,10 +28,22 @@ from sqlalchemy import select
from airflow.models.asset import AssetDagRunQueue, AssetEvent, AssetModel
from airflow.models.serialized_dag import SerializedDagModel
from airflow.providers.standard.operators.empty import EmptyOperator
-from airflow.sdk import Asset, AssetAll, AssetAny, AssetOrTimeSchedule as
SdkAssetOrTimeSchedule
+from airflow.sdk import (
+ Asset,
+ AssetAll,
+ AssetAndTimeSchedule as SdkAssetAndTimeSchedule,
+ AssetAny,
+ AssetOrTimeSchedule as SdkAssetOrTimeSchedule,
+)
+from airflow.sdk.bases.timetable import BaseTimetable
+from airflow.sdk.definitions.timetables.assets import AssetTriggeredTimetable
as SdkAssetTriggeredTimetable
+from airflow.sdk.exceptions import AirflowTimetableInvalid
from airflow.serialization.definitions.assets import SerializedAsset,
SerializedAssetAll, SerializedAssetAny
from airflow.serialization.serialized_objects import DagSerialization
-from airflow.timetables.assets import AssetOrTimeSchedule as
CoreAssetOrTimeSchedule
+from airflow.timetables.assets import (
+ AssetAndTimeSchedule as CoreAssetAndTimeSchedule,
+ AssetOrTimeSchedule as CoreAssetOrTimeSchedule,
+)
from airflow.timetables.base import DagRunInfo, DataInterval, TimeRestriction,
Timetable
from airflow.timetables.simple import AssetTriggeredTimetable
from airflow.utils.types import DagRunType
@@ -81,6 +93,14 @@ class MockTimetable(Timetable):
return DataInterval.exact(run_after)
+class CustomAssetTriggeredTimetable(MockTimetable):
+ asset_triggered = True
+
+
+class CustomAssetGatedTimetable(MockTimetable):
+ asset_gated = True
+
+
def serialize_timetable(timetable: Timetable) -> str:
"""
Mock serialization function for Timetable objects.
@@ -122,6 +142,19 @@ def sdk_asset_timetable(test_timetable, test_assets) ->
SdkAssetOrTimeSchedule:
return SdkAssetOrTimeSchedule(timetable=test_timetable, assets=test_assets)
[email protected]
+def sdk_asset_and_time_timetable(test_timetable, test_assets) ->
SdkAssetAndTimeSchedule:
+ return SdkAssetAndTimeSchedule(timetable=test_timetable,
assets=test_assets)
+
+
[email protected]
+def core_asset_and_time_timetable(test_timetable: MockTimetable) ->
CoreAssetAndTimeSchedule:
+ return CoreAssetAndTimeSchedule(
+ timetable=test_timetable,
+ assets=SerializedAssetAll([SerializedAsset("test_asset",
"test://asset/", "asset", {}, [])]),
+ )
+
+
@pytest.fixture
def core_asset_timetable(test_timetable: MockTimetable) ->
CoreAssetOrTimeSchedule:
return CoreAssetOrTimeSchedule(
@@ -130,6 +163,97 @@ def core_asset_timetable(test_timetable: MockTimetable) ->
CoreAssetOrTimeSchedu
)
[email protected](
+ ("timetable", "expected"),
+ [
+ pytest.param(MockTimetable(), (False, False), id="core-regular"),
+ pytest.param(
+ AssetTriggeredTimetable(SerializedAsset("test_asset",
"test://asset/", "asset", {}, [])),
+ (True, False),
+ id="core-asset-triggered",
+ ),
+ pytest.param(
+ CoreAssetOrTimeSchedule(
+ timetable=MockTimetable(),
+ assets=SerializedAsset("test_asset", "test://asset/", "asset",
{}, []),
+ ),
+ (True, False),
+ id="core-asset-or-time",
+ ),
+ pytest.param(
+ CoreAssetAndTimeSchedule(
+ timetable=MockTimetable(),
+ assets=SerializedAsset("test_asset", "test://asset/", "asset",
{}, []),
+ ),
+ (False, True),
+ id="core-asset-and-time",
+ ),
+ pytest.param(BaseTimetable(), (False, False), id="sdk-regular"),
+ pytest.param(
+ SdkAssetTriggeredTimetable(assets=Asset("test")),
+ (True, False),
+ id="sdk-asset-triggered",
+ ),
+ pytest.param(
+ SdkAssetOrTimeSchedule(timetable=BaseTimetable(),
assets=Asset("test")),
+ (True, False),
+ id="sdk-asset-or-time",
+ ),
+ pytest.param(
+ SdkAssetAndTimeSchedule(timetable=BaseTimetable(),
assets=Asset("test")),
+ (False, True),
+ id="sdk-asset-and-time",
+ ),
+ ],
+)
+def test_asset_scheduling_behavior_flags(timetable, expected) -> None:
+ assert (timetable.asset_triggered, timetable.asset_gated) == expected
+
+
[email protected](
+ "outer_type",
+ [CoreAssetOrTimeSchedule, CoreAssetAndTimeSchedule],
+)
[email protected](
+ "inner_type",
+ [CustomAssetTriggeredTimetable, CustomAssetGatedTimetable],
+)
+def
test_core_asset_time_schedules_reject_nested_asset_aware_timetable(outer_type,
inner_type) -> None:
+ asset = SerializedAsset("test_asset", "test://asset/", "asset", {}, [])
+ timetable = outer_type(timetable=inner_type(), assets=asset)
+
+ with pytest.raises(AirflowTimetableInvalid, match="Cannot nest asset-aware
timetables"):
+ timetable.validate()
+
+
[email protected](
+ ("assets", "expected_type"),
+ [
+ pytest.param(
+ SerializedAsset("test_asset", "test://asset/", "asset", {}, []),
+ SerializedAsset,
+ id="serialized-asset-used-as-is",
+ ),
+ pytest.param(
+ [SerializedAsset("test_asset", "test://asset/", "asset", {}, [])],
+ SerializedAssetAll,
+ id="collection-wrapped-in-all",
+ ),
+ pytest.param(Asset("test_asset"), SerializedAsset,
id="sdk-asset-converted"),
+ pytest.param([Asset("test_asset")], SerializedAssetAll,
id="sdk-collection-converted"),
+ ],
+)
+def test_core_asset_and_time_schedule_coerces_assets(assets, expected_type) ->
None:
+ timetable = CoreAssetAndTimeSchedule(timetable=MockTimetable(),
assets=assets)
+
+ assert isinstance(timetable.asset_condition, expected_type)
+
+
+def test_core_asset_and_time_schedule_rejects_non_asset() -> None:
+ with pytest.raises(ValueError, match="serialization not implemented for
'int'"):
+ CoreAssetAndTimeSchedule(timetable=MockTimetable(), assets=123) #
type: ignore[arg-type]
+
+
def test_serialization(sdk_asset_timetable: SdkAssetOrTimeSchedule,
monkeypatch: Any) -> None:
"""
Tests the serialization method of AssetOrTimeSchedule.
@@ -160,6 +284,31 @@ def test_serialization(sdk_asset_timetable:
SdkAssetOrTimeSchedule, monkeypatch:
}
+def test_serialization_and(sdk_asset_and_time_timetable:
SdkAssetAndTimeSchedule, monkeypatch: Any) -> None:
+ """Tests serialization of AssetAndTimeSchedule."""
+ from airflow.serialization.encoders import _serializer
+
+ monkeypatch.setattr(
+ "airflow.serialization.encoders.encode_timetable", lambda x:
"mock_serialized_timetable"
+ )
+ serialized = _serializer.serialize_timetable(sdk_asset_and_time_timetable)
+ assert serialized == {
+ "timetable": "mock_serialized_timetable",
+ "asset_condition": {
+ "__type": "asset_all",
+ "objects": [
+ {
+ "__type": "asset",
+ "name": "test_asset",
+ "uri": "test://asset/",
+ "group": "asset",
+ "extra": {},
+ }
+ ],
+ },
+ }
+
+
def test_deserialization(monkeypatch: Any, core_asset_timetable:
CoreAssetOrTimeSchedule) -> None:
"""
Tests the deserialization method of AssetOrTimeSchedule.
@@ -186,6 +335,37 @@ def test_deserialization(monkeypatch: Any,
core_asset_timetable: CoreAssetOrTime
assert deserialized == core_asset_timetable
+def test_deserialization_and(
+ monkeypatch: Any, core_asset_and_time_timetable: CoreAssetAndTimeSchedule
+) -> None:
+ """Tests deserialization of AssetAndTimeSchedule."""
+ monkeypatch.setattr("airflow.serialization.decoders.decode_timetable",
lambda x: MockTimetable())
+ mock_serialized_data = {
+ "timetable": "mock_serialized_timetable",
+ "asset_condition": {
+ "__type": "asset_all",
+ "objects": [
+ {
+ "__type": "asset",
+ "name": "test_asset",
+ "uri": "test://asset/",
+ "group": "asset",
+ "extra": None,
+ }
+ ],
+ },
+ }
+ deserialized = CoreAssetAndTimeSchedule.deserialize(mock_serialized_data)
+ assert isinstance(deserialized, CoreAssetAndTimeSchedule)
+ assert isinstance(deserialized.timetable, MockTimetable)
+ assert isinstance(deserialized.asset_condition, SerializedAssetAll)
+ assert len(deserialized.asset_condition.objects) == 1
+ asset = deserialized.asset_condition.objects[0]
+ assert isinstance(asset, SerializedAsset)
+ assert asset.name == "test_asset"
+ assert asset.uri == "test://asset/"
+
+
def test_infer_manual_data_interval(core_asset_timetable:
CoreAssetOrTimeSchedule) -> None:
"""
Tests the infer_manual_data_interval method of AssetOrTimeSchedule.
@@ -197,6 +377,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(2025, 6, 7, 8, 9, tzinfo=UTC)
+ result =
core_asset_and_time_timetable.infer_manual_data_interval(run_after=run_after)
+ assert result == DataInterval.exact(run_after)
+
+
def test_next_dagrun_info(core_asset_timetable: CoreAssetOrTimeSchedule) ->
None:
"""
Tests the next_dagrun_info method of AssetOrTimeSchedule.
@@ -214,6 +400,18 @@ def test_next_dagrun_info(core_asset_timetable:
CoreAssetOrTimeSchedule) -> None
)
+def test_next_dagrun_info_and(core_asset_and_time_timetable:
CoreAssetAndTimeSchedule) -> None:
+ last_interval = DataInterval.exact(DateTime(2025, 6, 7, 8, 9, tzinfo=UTC))
+ restriction = TimeRestriction(earliest=DateTime(2025, 6, 9, 8, 9,
tzinfo=UTC), latest=None, catchup=True)
+ result = core_asset_and_time_timetable.next_dagrun_info(
+ last_automated_data_interval=last_interval, restriction=restriction
+ )
+ assert result == DagRunInfo.interval(
+ DateTime(2025, 6, 9, 8, 9, tzinfo=UTC),
+ DateTime(2025, 6, 10, 8, 9, tzinfo=UTC),
+ )
+
+
def test_generate_run_id(core_asset_timetable: CoreAssetOrTimeSchedule) ->
None:
"""
Tests the generate_run_id method of AssetOrTimeSchedule.
@@ -231,6 +429,31 @@ def test_generate_run_id(core_asset_timetable:
CoreAssetOrTimeSchedule) -> None:
assert run_id == "manual__2025-06-07T08:09:00+00:00"
+def test_generate_run_id_and(core_asset_and_time_timetable:
CoreAssetAndTimeSchedule, mocker) -> None:
+ date = DateTime(2025, 6, 7, 8, 9, tzinfo=UTC)
+ generate_run_id = mocker.patch.object(
+ core_asset_and_time_timetable.timetable,
+ "generate_run_id",
+ autospec=True,
+ return_value="wrapped_run_id",
+ )
+ run_id = core_asset_and_time_timetable.generate_run_id(
+ run_type=DagRunType.MANUAL,
+ extra_args="test",
+ logical_date=date,
+ run_after=date,
+ data_interval=None,
+ )
+ assert run_id == "wrapped_run_id"
+ generate_run_id.assert_called_once_with(
+ run_type=DagRunType.MANUAL,
+ extra_args="test",
+ logical_date=date,
+ run_after=date,
+ data_interval=None,
+ )
+
+
@pytest.fixture
def asset_events(mocker) -> list[AssetEvent]:
"""Pytest fixture for creating mock AssetEvent objects."""
diff --git a/task-sdk/docs/api.rst b/task-sdk/docs/api.rst
index 6ce6213e0c3..4f83dc1d838 100644
--- a/task-sdk/docs/api.rst
+++ b/task-sdk/docs/api.rst
@@ -204,6 +204,8 @@ Assets
Timetables
----------
+.. autoapiclass:: airflow.sdk.AssetAndTimeSchedule
+
.. autoapiclass:: airflow.sdk.AssetOrTimeSchedule
.. autoapiclass:: airflow.sdk.CronDataIntervalTimetable
diff --git a/task-sdk/src/airflow/sdk/__init__.py
b/task-sdk/src/airflow/sdk/__init__.py
index aec71b22290..91822202553 100644
--- a/task-sdk/src/airflow/sdk/__init__.py
+++ b/task-sdk/src/airflow/sdk/__init__.py
@@ -26,6 +26,7 @@ __all__ = [
"AssetAlias",
"AssetAll",
"AssetAny",
+ "AssetAndTimeSchedule",
"AssetOrTimeSchedule",
"AssetWatcher",
"AsyncCallback",
@@ -207,6 +208,7 @@ if TYPE_CHECKING:
from airflow.sdk.definitions.taskgroup import TaskGroup
from airflow.sdk.definitions.template import literal
from airflow.sdk.definitions.timetables.assets import (
+ AssetAndTimeSchedule,
AssetOrTimeSchedule,
PartitionedAssetTimetable,
PartitionedAtRuntime,
@@ -237,6 +239,7 @@ __lazy_imports: dict[str, str] = {
"AssetAccessControl": ".definitions.asset",
"AssetAlias": ".definitions.asset",
"AssetAll": ".definitions.asset",
+ "AssetAndTimeSchedule": ".definitions.timetables.assets",
"AssetAny": ".definitions.asset",
"AssetOrTimeSchedule": ".definitions.timetables.assets",
"AssetWatcher": ".definitions.asset",
diff --git a/task-sdk/src/airflow/sdk/__init__.pyi
b/task-sdk/src/airflow/sdk/__init__.pyi
index 6e7a0c8526c..c3a94d30a35 100644
--- a/task-sdk/src/airflow/sdk/__init__.pyi
+++ b/task-sdk/src/airflow/sdk/__init__.pyi
@@ -112,6 +112,7 @@ from airflow.sdk.definitions.retry_policy import (
from airflow.sdk.definitions.taskgroup import TaskGroup as TaskGroup
from airflow.sdk.definitions.template import literal as literal
from airflow.sdk.definitions.timetables.assets import (
+ AssetAndTimeSchedule,
AssetOrTimeSchedule,
PartitionedAssetTimetable,
PartitionedAtRuntime,
@@ -143,6 +144,7 @@ __all__ = [
"AssetAccessControl",
"AssetAlias",
"AssetAll",
+ "AssetAndTimeSchedule",
"AssetAny",
"AssetOrTimeSchedule",
"AssetWatcher",
diff --git a/task-sdk/src/airflow/sdk/bases/timetable.py
b/task-sdk/src/airflow/sdk/bases/timetable.py
index 27769d6b31c..735d8b35e00 100644
--- a/task-sdk/src/airflow/sdk/bases/timetable.py
+++ b/task-sdk/src/airflow/sdk/bases/timetable.py
@@ -47,6 +47,13 @@ class BaseTimetable:
asset_condition: BaseAsset | None = None
+ # TODO (GH-52141): Find a way to keep these and the ones in Core in sync.
+ asset_triggered: bool = False
+ """Whether this timetable creates runs triggered by asset events."""
+
+ asset_gated: bool = False
+ """Whether this timetable's scheduled runs are gated on an asset
condition."""
+
partitioned_at_runtime: bool = False
"""
Whether this timetable defers partition selection to task runtime.
diff --git a/task-sdk/src/airflow/sdk/definitions/dag.py
b/task-sdk/src/airflow/sdk/definitions/dag.py
index de1cabb34bc..c3d485eb91b 100644
--- a/task-sdk/src/airflow/sdk/definitions/dag.py
+++ b/task-sdk/src/airflow/sdk/definitions/dag.py
@@ -640,7 +640,7 @@ class DAG:
return
from airflow.sdk.api.datamodels._generated import DagRunType
- if isinstance(self.timetable, AssetTriggeredTimetable):
+ if self.timetable.asset_triggered:
if DagRunType.ASSET_TRIGGERED not in allowed_run_types:
raise ValueError(
"allowed_run_types must include ASSET_TRIGGERED when the
Dag is scheduled by assets"
diff --git a/task-sdk/src/airflow/sdk/definitions/timetables/assets.py
b/task-sdk/src/airflow/sdk/definitions/timetables/assets.py
index 22208693588..bc09cd8f2b4 100644
--- a/task-sdk/src/airflow/sdk/definitions/timetables/assets.py
+++ b/task-sdk/src/airflow/sdk/definitions/timetables/assets.py
@@ -42,6 +42,7 @@ class AssetTriggeredTimetable(BaseTimetable):
:meta private:
"""
+ asset_triggered = True
asset_condition: BaseAsset = attrs.field(alias="assets")
@@ -84,3 +85,27 @@ class AssetOrTimeSchedule(AssetTriggeredTimetable):
def __attrs_post_init__(self) -> None:
self.active_runs_limit = self.timetable.active_runs_limit
self.can_be_scheduled = self.timetable.can_be_scheduled
+
+
[email protected](kw_only=True)
+class AssetAndTimeSchedule(BaseTimetable):
+ """
+ Combine time-based scheduling with asset conditions.
+
+ :param assets: An asset or list of assets, in the same format as
+ ``DAG(schedule=...)`` when using event-driven scheduling. This is used
+ to evaluate whether a scheduled run can be created.
+ :param timetable: A timetable instance to evaluate time-based scheduling.
+ """
+
+ asset_gated = True
+ asset_condition: BaseAsset = attrs.field(alias="assets",
converter=_coerce_assets)
+ timetable: BaseTimetable
+
+ @property
+ def active_runs_limit(self) -> int | None: # type: ignore[override]
+ return self.timetable.active_runs_limit
+
+ @property
+ def can_be_scheduled(self) -> bool: # type: ignore[override]
+ return self.timetable.can_be_scheduled
diff --git a/task-sdk/tests/task_sdk/definitions/test_dag.py
b/task-sdk/tests/task_sdk/definitions/test_dag.py
index 07b8c8186c8..0925ad2df85 100644
--- a/task-sdk/tests/task_sdk/definitions/test_dag.py
+++ b/task-sdk/tests/task_sdk/definitions/test_dag.py
@@ -27,6 +27,7 @@ import pytest
from airflow.sdk import (
DAG,
+ Asset,
Context,
Label,
Param,
@@ -638,8 +639,6 @@ def test_allowed_run_types_conflicting_schedule(schedule,
allowed_run_types, mat
def test_allowed_run_types_asset_triggered_missing_with_asset_schedule():
- from airflow.sdk.definitions.asset import Asset
-
with pytest.raises(ValueError, match="allowed_run_types must include
ASSET_TRIGGERED"):
DAG(
"test-allowed-asset",
@@ -648,6 +647,19 @@ def
test_allowed_run_types_asset_triggered_missing_with_asset_schedule():
)
+def test_allowed_run_types_uses_asset_triggered_behavior():
+ class CustomAssetTriggeredTimetable(BaseTimetable):
+ asset_triggered = True
+ asset_condition = Asset("test")
+
+ with pytest.raises(ValueError, match="allowed_run_types must include
ASSET_TRIGGERED"):
+ DAG(
+ "test-allowed-custom-asset",
+ schedule=CustomAssetTriggeredTimetable(),
+ allowed_run_types=[DagRunType.MANUAL],
+ )
+
+
def test__tags_mutable():
expected_tags = {"6", "7"}
test_dag = DAG("test-dag")