Pebble32 commented on code in PR #72475:
URL: https://github.com/apache/airflow/pull/72475#discussion_r4125246239
##########
airflow-core/src/airflow/timetables/trigger.py:
##########
@@ -406,8 +441,12 @@ def __init__(
run_offset: int | datetime.timedelta | relativedelta | None = None,
run_immediately: bool | datetime.timedelta = False,
key_format: str = r"%Y-%m-%dT%H:%M:%S",
+ seed: str = "",
+ max_jitter: datetime.timedelta = datetime.timedelta(),
Review Comment:
Fixed, keys now stay on the cron boundary while the run fires late. Added
tests for both of your examples (including the coarse `key_format` date flip)
and for backfill iteration with `run_offset` 0, 1 and -1. ✅
##########
airflow-core/src/airflow/timetables/_cron.py:
##########
@@ -61,15 +62,54 @@ def _covers_every_hour(cron: croniter) -> bool:
class CronMixin:
- """Mixin to provide interface to work with croniter."""
+ """
+ Mixin to provide interface to work with croniter.
+
+ Optionally applies a deterministic, per-DAG jitter to every scheduled time.
+ When ``max_jitter`` is set, each cron boundary is shifted by a fixed offset
+ derived from ``seed`` and spread across ``[0, max_jitter)``. This spreads
out
+ DAGs that share a cron expression (e.g. every ``@daily`` DAG firing at
+ midnight) so they no longer all fire at the same instant. The offset is
stable
+ for a given seed, so runs stay predictable across scheduler restarts and
+ serialization.
+
+ The offset shifts every cron-derived time uniformly. For data-interval
+ timetables this means the whole interval moves by the offset (the window
keeps
+ its length and consecutive intervals stay contiguous), not just the fire
time.
+
+ :param cron: cron expression (or a preset such as ``@daily``) defining the
schedule.
+ :param timezone: timezone used to interpret the cron expression.
+ :param seed: stable, unique-per-DAG string the offset is derived from; the
DAG id
+ is a natural choice. Must be non-empty whenever ``max_jitter`` is set.
+ :param max_jitter: upper bound of the jitter window; the offset falls in
+ ``[0, max_jitter)``. Defaults to zero, i.e. no jitter. Keep it small
relative
+ to the gap between cron boundaries.
+ """
- def __init__(self, cron: str, timezone: str | Timezone | FixedTimezone) ->
None:
+ def __init__(
+ self,
+ cron: str,
+ timezone: str | Timezone | FixedTimezone,
+ *,
+ seed: str = "",
+ max_jitter: datetime.timedelta = datetime.timedelta(),
+ ) -> None:
self._expression = cron_presets.get(cron, cron)
if isinstance(timezone, str):
timezone = parse_timezone(timezone)
self._timezone = timezone
+ if max_jitter > datetime.timedelta(0) and not seed:
+ raise ValueError("seed must be a non-empty, unique-per-DAG string
when max_jitter > 0")
Review Comment:
Fixed, now rejected in both layers. Added tests for core and SDK. ✅
##########
airflow-core/src/airflow/timetables/_cron.py:
##########
@@ -61,15 +62,54 @@ def _covers_every_hour(cron: croniter) -> bool:
class CronMixin:
- """Mixin to provide interface to work with croniter."""
+ """
+ Mixin to provide interface to work with croniter.
+
+ Optionally applies a deterministic, per-DAG jitter to every scheduled time.
+ When ``max_jitter`` is set, each cron boundary is shifted by a fixed offset
+ derived from ``seed`` and spread across ``[0, max_jitter)``. This spreads
out
+ DAGs that share a cron expression (e.g. every ``@daily`` DAG firing at
+ midnight) so they no longer all fire at the same instant. The offset is
stable
+ for a given seed, so runs stay predictable across scheduler restarts and
+ serialization.
+
+ The offset shifts every cron-derived time uniformly. For data-interval
+ timetables this means the whole interval moves by the offset (the window
keeps
+ its length and consecutive intervals stay contiguous), not just the fire
time.
+
+ :param cron: cron expression (or a preset such as ``@daily``) defining the
schedule.
+ :param timezone: timezone used to interpret the cron expression.
+ :param seed: stable, unique-per-DAG string the offset is derived from; the
DAG id
+ is a natural choice. Must be non-empty whenever ``max_jitter`` is set.
+ :param max_jitter: upper bound of the jitter window; the offset falls in
+ ``[0, max_jitter)``. Defaults to zero, i.e. no jitter. Keep it small
relative
+ to the gap between cron boundaries.
+ """
- def __init__(self, cron: str, timezone: str | Timezone | FixedTimezone) ->
None:
+ def __init__(
+ self,
+ cron: str,
+ timezone: str | Timezone | FixedTimezone,
+ *,
+ seed: str = "",
+ max_jitter: datetime.timedelta = datetime.timedelta(),
+ ) -> None:
self._expression = cron_presets.get(cron, cron)
if isinstance(timezone, str):
timezone = parse_timezone(timezone)
self._timezone = timezone
+ if max_jitter > datetime.timedelta(0) and not seed:
+ raise ValueError("seed must be a non-empty, unique-per-DAG string
when max_jitter > 0")
+ h = int(md5(seed.encode()).hexdigest(), 16)
Review Comment:
Moved inside the jitter branch. ✅
##########
airflow-core/src/airflow/timetables/_cron.py:
##########
@@ -61,15 +62,54 @@ def _covers_every_hour(cron: croniter) -> bool:
class CronMixin:
- """Mixin to provide interface to work with croniter."""
+ """
+ Mixin to provide interface to work with croniter.
+
+ Optionally applies a deterministic, per-DAG jitter to every scheduled time.
Review Comment:
Fixed on all lines this PR adds. ✅
##########
airflow-core/src/airflow/timetables/_cron.py:
##########
@@ -120,18 +168,26 @@ def _describe_with_dom_dow_fix(self, expression: str) ->
str:
def __eq__(self, other: object) -> bool:
"""
- Both expression and timezone should match.
+ Expression, timezone and jitter settings (``seed`` and ``max_jitter``)
should all match.
+
+ Two timetables that share a cron expression and timezone but differ in
jitter
+ produce different schedules, so they are not considered equal.
This is only for testing purposes and should not be relied on
otherwise.
"""
from airflow.serialization.encoders import coerce_to_core_timetable
if not isinstance(other := coerce_to_core_timetable(other),
type(self)):
return NotImplemented
- return self._expression == other._expression and self._timezone ==
other._timezone
+ return (
+ self._expression == other._expression
+ and self._timezone == other._timezone
+ and self._seed == other._seed
+ and self._max_jitter == other._max_jitter
+ )
def __hash__(self):
- return hash((self._expression, self._timezone))
+ return hash((self._expression, str(self._timezone), self._seed,
self._max_jitter))
Review Comment:
Split out as #73859 with its own test. ✅
##########
task-sdk/src/airflow/sdk/definitions/timetables/_cron.py:
##########
@@ -41,16 +42,33 @@
@attrs.define
class CronMixin:
- """Mixin to provide interface to work with croniter."""
+ """
+ Mixin to provide interface to work with croniter.
+
+ Optionally applies a deterministic, per-DAG jitter: when ``max_jitter`` is
set, every
+ cron boundary is shifted by a fixed offset derived from ``seed`` and
spread across
+ ``[0, max_jitter)``, so DAGs sharing a cron expression no longer all fire
at the same
+ instant. The offset is computed scheduler-side; this class only carries
and validates
+ the settings.
+
+ :param seed: stable, unique-per-DAG string the offset is derived from (the
DAG id is a
+ natural choice). Must be non-empty whenever ``max_jitter`` is set.
+ :param max_jitter: upper bound of the jitter window; the offset falls in
+ ``[0, max_jitter)``. Defaults to zero, i.e. no jitter.
+ """
expression: str
timezone: str | Timezone | FixedTimezone
+ seed: str = attrs.field(kw_only=True, default="")
+ max_jitter: datetime.timedelta = attrs.field(kw_only=True,
default=datetime.timedelta())
def __attrs_post_init__(self) -> None:
# Resolve preset aliases (e.g. "@quarterly") to their cron expressions
# in-place. After this point the original preset string is lost;
# attrs.evolve, equality, and serialisation all see the resolved form.
self.expression = CRON_PRESETS.get(self.expression, self.expression)
+ if self.max_jitter > datetime.timedelta(0) and not self.seed:
+ raise ValueError("seed must be a non-empty, unique-per-DAG string
when max_jitter > 0")
Review Comment:
Added to `test__cron.py`: empty seed, negative window, and defaults being a
no-op. ✅
--
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]