Repository navigation
Conversation
|
Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contributors' Guide
|
a11aa56 to
5f777d1
Compare
Signed-off-by: Adam <111773160+Pebble32@users.noreply.github.com>
Signed-off-by: Adam <111773160+Pebble32@users.noreply.github.com>
Signed-off-by: Adam <111773160+Pebble32@users.noreply.github.com>
5f777d1 to
b668d8d
Compare
…ed guard Signed-off-by: Adam <111773160+Pebble32@users.noreply.github.com>
Signed-off-by: Adam <111773160+Pebble32@users.noreply.github.com>
|
I wonder if it’d make sense to add this functionality directly to CronMixin (and therefore inherited by all cron-based scheduling). Some scheduling tools (Jenkins IIRC) have this built-in for cron scheduling (opt-in) to not let a large number of jobs spiking a server periodically if timing is not essential. Implementation seems reasonable to me in general. |
|
Thanks @uranusjr! The Jenkins H comparison is exactly the idea I kept it as its own timetable to keep the change contained. Moving it into Would you prefer I land this focused version first and generalize into |
|
hey about "enough to cause task failures" if it can happen then it will happen again ( backfill , new dags , big clear ... ) your tasks fail because you did not put correct/perfect limit on concurrency. I know it's not easy : yes airflow pools are too simple for many use-case yes airflow is not kubernetes ressource aware ( he do not know your max hardware scaling limits and don't queue tasks regarding this ) "but never spread the fire times apart at the source." I've a stack triggering more than a thousand of dag_run at midnight , yes it take almost 30 seconds to the scheduler to do so , but no errors , did you encounter a scheduler error ? |
|
Fair point, and well put. Where I still think it earns its place is as peak shaving on top of concurrency limits, not instead of them: spreading arrivals so fewer tasks land in the pool at the same instant, so a fixed pool saturates less often. Opt in, for when exact timing does not matter. Same idea as Jenkins H that @uranusjr mentioned. Maybe it is better to move it directly to CronMixin |
|
Can you explore this locally to see how intrusive adding this would be? If implementation becomes messy, I think it’s reasonable to do this in a separate timetable in this PR first with the intention to eventually refactor the logic into CronMixin before 3.4.0 is released, which is quite still some time away. |
|
Sounds good. I will explore the |
|
Explored it and it seems less intrusive than I first thought. The offset moves into For CronDataIntervalTimetable the shared offset shifts both interval bounds and the fire time by the same amount, so the window stays one full period long but sits offset from the cron line (00:35 to 00:35 instead of 00:00 to 00:00), and consecutive runs stay contiguous so no gaps or overlaps. I will go with that uniform shift as the default and document it, keeping the guidance to set max_jitter small relative to the period. Jittering only the fire time while keeping the interval on the cron boundary is possible but needs a targeted override, so I will leave it as a possible follow up. Will put up the CronMixin version. |
|
Follow up moving the jitter into CronMixin as discussed: #72475. Will close this PR once that one lands. |
Add
JitteredCronTimetable, an opt-inCronTriggerTimetablesubclass that shifts each DAG's fire time by a deterministic, per-DAG offset drawn from[0, max_jitter). This spreads out DAGs that share a cron expression so they no longer all fire at the same instant, without changinglogical_date/data_intervalsemantics.Why
@dailyexpands to0 0 * * *, so every daily DAG in a deployment is scheduled at exactly midnight. In large deployments this "thundering herd" at the cron boundary overloads the scheduler and workers — enough to cause task failures when dozens of DAGs are born at the same instant.The existing ways to deal with this don't actually de-collide the schedule:
@dailyintentparallelism,max_active_tasks_per_dag, pool slots)JitteredCronTimetableis the only approach that moves the fire times themselves, deterministically: the sameseed(e.g. the DAG id) always maps to the same offset, so runs stay stable and predictable across scheduler restarts and timetable serialization. Motivated by the discussion in #69027.What
JitteredCronTimetable(CronTriggerTimetable)in both the Task SDK (author-facing, attrs-based) and airflow-core (scheduler-side), plus the serialization wiring (BUILTIN_TIMETABLESmapping,serialize/deserialize, encode/decode across the SDK↔core boundary).CronTriggerTimetable:seed: strandmax_jitter: timedelta.md5(seed) % max_jitter.total_seconds(), applied as a "strip → cron → apply" coordinate shift so cron/DST alignment is fully delegated to the parent and only the wall-clock fire time is shifted.seed="",max_jitter=timedelta(0)) the offset is zero and it behaves identically toCronTriggerTimetable. Nothing changes for anyone who doesn't use it.Tests
airflow-core/tests/unit/timetables/test_jittered_cron_timetable.py, modeled on the existingtest_trigger_timetable.py:America/New_York);max_jitterreproducesCronTriggerTimetableexactly (catchup on/off);[0, max_jitter), and spread across distinct seeds;serialize/deserializeround-trips seed + window + derived offset;Was generative AI tooling used to co-author this PR?
Generated-by: Claude (Claude Code), following the guidelines
related: #69027