GitHub user raphaelauv added a comment to the discussion: Manual DAG execution
with defined delay
best solution would be with TimeDeltaSensor , but the argument is not templated
( dynamic )
so only solution at the moment is doing
```python
from time import sleep
from typing import TYPE_CHECKING, Sequence
from airflow.providers.common.compat.sdk import BaseSensorOperator, conf
from airflow.providers.standard.triggers.temporal import TimeDeltaTrigger
if TYPE_CHECKING:
from airflow.providers.common.compat.sdk import Context
class MyWaitSensor(BaseSensorOperator):
template_fields: Sequence[str] = ("time_to_wait",)
def __init__(
self,
time_to_wait: timedelta | int | str,
deferrable: bool = conf.getboolean("operators",
"default_deferrable", fallback=False),
**kwargs,
) -> None:
super().__init__(**kwargs)
self.deferrable = deferrable
self.time_to_wait = time_to_wait
def _resolve_time_to_wait(self) -> timedelta:
value = self.time_to_wait
if isinstance(value, timedelta):
return value
else:
return timedelta(seconds=int(value))
def execute(self, context: Context) -> None:
time_to_wait = self._resolve_time_to_wait()
if self.deferrable:
self.defer(
trigger=(
TimeDeltaTrigger(time_to_wait, end_from_trigger=True)
),
method_name="execute_complete",
)
else:
sleep(int(time_to_wait.total_seconds()))
from airflow.sdk import DAG, Param
with DAG(
dag_id="my_dag",
start_date=datetime(2024, 1, 1),
schedule=None,
catchup=False,
params={
"duration_seconds": Param(
default=30,
type="integer",
minimum=1,
description="Duration in seconds to wait before continuing.",
)
},
):
MyWaitSensor(
task_id="wait_for_duration",
time_to_wait="{{ params.duration_seconds }}",
deferrable=True
)
```
GitHub link:
https://github.com/apache/airflow/discussions/68383#discussioncomment-17790345
----
This is an automatically sent email for [email protected].
To unsubscribe, please send an email to: [email protected]