Lee-W commented on code in PR #58543:
URL: https://github.com/apache/airflow/pull/58543#discussion_r4141914252
##########
airflow-core/src/airflow/models/dag.py:
##########
@@ -131,8 +132,10 @@ def infer_automated_data_interval(timetable: Timetable,
logical_date: datetime)
:meta private:
"""
+ if timetable.asset_triggered:
+ return DataInterval.exact(timezone.coerce_datetime(logical_date))
timetable_type = type(timetable)
- if issubclass(timetable_type, (NullTimetable, OnceTimetable,
AssetTriggeredTimetable)):
+ if issubclass(timetable_type, (NullTimetable, OnceTimetable)):
return DataInterval.exact(timezone.coerce_datetime(logical_date))
Review Comment:
```suggestion
if (
timetable.asset_triggered
or issubclass(timetable_type := type(timetable), (NullTimetable,
OnceTimetable))
):
return
DataInterval.exact(timezone.coerce_datetime(logical_date))
```
##########
airflow-core/tests/unit/timetables/test_assets_timetable.py:
##########
@@ -130,6 +163,114 @@ 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", "inner_type"),
+ [
+ pytest.param(
+ CoreAssetOrTimeSchedule,
+ CustomAssetTriggeredTimetable,
+ id="or-wraps-custom-asset-triggered",
+ ),
+ pytest.param(
+ CoreAssetOrTimeSchedule,
+ CustomAssetGatedTimetable,
+ id="or-wraps-custom-asset-gated",
+ ),
+ pytest.param(
+ CoreAssetAndTimeSchedule,
+ CustomAssetTriggeredTimetable,
+ id="and-wraps-custom-asset-triggered",
+ ),
+ pytest.param(
+ CoreAssetAndTimeSchedule,
+ CustomAssetGatedTimetable,
+ id="and-wraps-custom-asset-gated",
+ ),
+ ],
+)
Review Comment:
```suggestion
@pytest.mark.parametrize(
"outer_type",
[CoreAssetOrTimeSchedule, CoreAssetAndTimeSchedule],
)
@pytest.mark.parametrize(
"inner_type",
[CustomAssetTriggeredTimetable, CustomAssetGatedTimetable],
)
```
##########
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")
Review Comment:
```suggestion
raise AirflowTimetableInvalid("Cannot nest asset-aware timetables")
if not isinstance(asset_condition, SerializedAssetBase):
raise AirflowTimetableInvalid("All elements in 'assets' must be
assets")
```
--
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]