Skip to content

Add opt-in deterministic jitter to cron-based timetables - #72475

Open
Pebble32 wants to merge 5 commits into
apache:mainfrom
Pebble32:add-cron-timetable-jitter
Open

Pebble32 wants to merge 5 commits into
apache:mainfrom
Pebble32:add-cron-timetable-jitter

Conversation

@Pebble32

@Pebble32 Pebble32 commented Sep 3, 2026

Copy link
Copy Markdown

Add opt-in, deterministic jitter to every cron-based timetable via CronMixin: two kw-only params, seed and max_jitter, on CronTriggerTimetable, CronDataIntervalTimetable, MultipleCronTriggerTimetable and CronPartitionTimetable, in both the Task SDK and airflow-core, with serialization wired through. Each DAG's runs are shifted by a fixed offset drawn from [0, max_jitter), so DAGs that share a cron expression no longer all fire at the same instant.

This is the follow-up @uranusjr suggested on #69705 (#69705 (comment)): rather than a standalone timetable, the jitter lives in CronMixin so all cron scheduling inherits it. It supersedes #69705, which I will close once this lands. cc @kaxil, who reviewed the original.

Why

@daily expands to 0 0 * * *, so every daily DAG in a deployment is scheduled at exactly midnight, and the same holds for any shared cron expression. The existing ways to deal with this don't actually spread the schedule:

Existing option Why it doesn't spread the schedule
Hand-pick a unique minute per DAG Manual, drifts and re-collides as the DAG count grows, and throws away the @daily intent
Hash the DAG id into a literal cron string Same loss of intent; not reusable across DAGs/teams; every author re-implements it ad hoc
Pools / concurrency limits (parallelism, max_active_tasks_per_dag, pool slots) Cap how many tasks run at once, but the runs are still scheduled at the same instant. They manage contention downstream; jitter reduces the peak at the source. The two are complementary, not alternatives.

Jitter is peak shaving on top of concurrency control, not a replacement for it. It is the same idea as Jenkins' H cron syntax. The offset is deterministic: the same seed (e.g. the DAG id) always maps to the same offset, so runs stay stable across scheduler restarts and serialization. Motivated by the discussion in #69027.

What

  • seed/max_jitter added to CronMixin in both layers, so all four cron timetables inherit them; the concrete classes thread the params through their constructors and serialize/deserialize, and the encoders.py variants emit them.
  • The offset is md5(seed) % max_jitter, computed in integer microseconds (FIPS-safe hashlib_wrapper.md5), and applied as a "strip → cron → apply" coordinate shift in CronMixin._get_next/_get_prev. Because _align_to_next/_align_to_prev build on those primitives, every cron timetable inherits the shift with no per-class scheduling code.
  • Fully opt-in and safe by default: with max_jitter at its default of zero the offset is zero and every timetable behaves identically to before. Nothing changes for anyone who doesn't use it.

Design decisions

  • Uniform shift for data-interval timetables. The cron boundaries define the window, so the whole interval moves by the offset: same length, consecutive runs stay contiguous, but it no longer starts exactly on the cron time (00:35 to 00:35 instead of 00:00 to 00:00). Documented, with guidance to keep max_jitter small relative to the period.
  • MultipleCronTriggerTimetable children share one offset. The same seed/max_jitter is passed to every child, so the whole DAG shifts in lockstep and its fire times keep their relative spacing.
  • Serialization only emits seed/max_jitter when jitter is set. The wire format of existing DAGs is unchanged (so nothing is re-serialized on upgrade) and deserialize tolerates missing keys. The existing serialization tests pass untouched.
  • Jitter is part of __eq__/__hash__, since timetables differing only in jitter produce different schedules. While adding this I fixed a latent bug: __hash__ included the pendulum Timezone object, which is unhashable, so hash() on any cron timetable raised; it now hashes str(timezone).
  • Empty seed with max_jitter > 0 raises in both layers: it would give every DAG the same offset and merely move the herd instead of spreading it.

Tests

airflow-core/tests/unit/timetables/test_cron_timetable_jitter.py (28 tests), modeled on the existing test_trigger_timetable.py:

  • jittered runs equal plain cron runs shifted by the fixed offset, across a catchup sequence including a DST spring-forward (America/New_York);
  • zero max_jitter reproduces the plain timetable exactly (catchup on/off);
  • offsets are deterministic for a given seed, bounded to [0, max_jitter), spread across distinct seeds, and correct for sub-second windows;
  • the empty-seed guard fires in both layers for all four timetables;
  • core serialize/deserialize and SDK encode → decode round-trips preserve the offset for all four timetables;
  • data-interval jitter shifts both bounds uniformly and keeps consecutive runs contiguous;
  • MultipleCronTriggerTimetable children share a single offset;
  • un-jittered timetables serialize without the jitter keys;
  • jitter participates in equality, and equal timetables hash equal.

Was generative AI tooling used to co-author this PR?
  • Yes (please specify the tool below)

Generated-by: Claude (Claude Code), following the guidelines


related: #69705
related: #69027

@potiuk potiuk left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for this — the design is clean (jitter in CronMixin so all four cron timetables inherit it), the opt-in default is genuinely a no-op, and the test file is thorough: DST-crossing catchup, sub-second windows, cross-layer encode/decode, and the "un-jittered Dags serialize unchanged" case are all exactly the right things to pin down.

One blocking issue, plus some smaller things inline.

Blocking

CronPartitionTimetable partition keys absorb the jitter.

_get_partition_date() returns run_date as-is for run_offset == 0, and _format_key() formats that — so now that _get_next/_align_to_next are jittered, the partition key shifts with the fire time. Using this PR's own offset function:

seed="my_dag", max_jitter=1h  -> offset 0:58:51.663322
  cron "0 0 * * *"            -> key 2026-03-06T00:58:51  (was 2026-03-06T00:00:00)

seed="dag_4",  max_jitter=2h  -> offset 1:03:38.886020
  cron "0 23 * * *", key_format="%Y-%m-%d"
                              -> key 2026-03-07           (was 2026-03-06)

Jitter is meant to move when a run executes, not which period it is for. The partition key is the run's identity — iter_partition_dagrun_infos dedups backfills on it and it typically maps to a storage path — so enabling jitter on an existing Dag renumbers every partition and stops it deduping against past ones. The second example shows the calendar date itself flipping for a coarser key_format.

I think the offset should be stripped before the partition date is derived, so the key stays on the cron boundary while the run still fires late. If the shifted key is deliberate, it needs to be documented in the new docs section and covered by a test — right now the CronPartitionTimetable tests only assert the serialize round-trip.

Smaller observations

See the inline comments: negative max_jitter is accepted and serialized, the md5 runs on every construction even with jitter off, the new docstrings use DAG where AGENTS.md asks for Dag (the .rst in this PR gets it right), and the SDK-side empty-seed guard is only tested from airflow-core/tests/ rather than task-sdk/tests/.../test__cron.py.

Two things that aren't findings:

  • summary / description don't mention the jitter. A jittered Dag renders as "0 0 * * * / At 12:00 AM" in the UI while actually firing at 00:58. Worth deciding whether summary should surface the offset — otherwise the schedule shown and the schedule run disagree, which is a support question waiting to happen.
  • seed is hand-written per Dag. The natural seed is the dag_id, but the timetable cannot see it. Copying a Dag file — a very common pattern — silently reuses the seed and recreates exactly the collision this feature exists to prevent. Is there room to default seed to the dag_id when the timetable is attached to a Dag, instead of raising on an empty seed?

The __hash__ fix is a real one — I confirmed hash() raises TypeError: unhashable type: 'Timezone' on current main for any cron timetable. Since it is an independent bugfix that is backportable on its own, consider splitting it into its own PR rather than having it ride along with a feature.

CI here is green but last ran on 2026-09-04. I test-merged this onto current main locally and it is clean and coherent, but please rebase for a fresh run before this goes in.

Worth a second look from

This touches the timetable core and both layers of serialization:

  • @uranusjr — suggested this shape in #69705 and owns timetables/
  • @kaxil — reviewed the original #69705

Neither has been notified — asking them is the maintainer's call, and optional.


This review was drafted by an AI-assisted tool and confirmed by an Apache Airflow maintainer. After you've addressed the points above and pushed an update, an Apache Airflow maintainer — a real person — will take the next look at the PR. The findings cite the project's review criteria; if you think one of them is mis-applied, please reply on the PR and a maintainer will weigh in.

More on how Apache Airflow handles maintainer review: contributing-docs/05_pull_requests.rst.


Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting

Comment thread airflow-core/src/airflow/timetables/trigger.py
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")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

minor — a negative max_jitter slips through both guards: this check is max_jitter > timedelta(0), and the offset branch below is if max_jitter_us > 0, so the offset ends up zero. But if self._max_jitter: in serialize() is truthy for a negative timedelta, so the payload still gets "max_jitter": -3600.0 and "seed": "". Silently doing nothing while persisting a nonsensical value — worth rejecting max_jitter < 0 outright here.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed, now rejected in both layers. Added tests for core and SDK. ✅


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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit — this runs on every CronMixin construction, including the default no-jitter path, i.e. for every cron Dag on every deserialization in the scheduler. It is cheap, but it only matters when max_jitter_us > 0, so it could move inside the branch below.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Moved inside the jitter branch. ✅

"""
Mixin to provide interface to work with croniter.

Optionally applies a deterministic, per-DAG jitter to every scheduled time.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

minor — AGENTS.md asks for Dag in prose:

Write Dag (title case) in all prose. Keep the all-caps or lowercase spelling only when reproducing a literal code token

This docstring has "per-DAG jitter", "every @daily DAG", "the DAG id is a natural choice"; the same applies to the :param seed: docs in trigger.py, both task-sdk files, the test docstrings, and the ValueError message ("unique-per-DAG string"). The .rst section added in this PR already uses Dag, so it is just inconsistent within the PR.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed on all lines this PR adds. ✅


def __hash__(self):
return hash((self._expression, self._timezone))
return hash((self._expression, str(self._timezone), self._seed, self._max_jitter))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Confirmed this is a real latent bug, not a theoretical one — on current main, hash(CronTriggerTimetable("0 0 * * *", timezone="UTC")) raises TypeError: unhashable type: 'Timezone'. Good catch.

Since it is independent of the jitter feature and backportable on its own, consider pulling it into a separate PR so it can land without waiting on this one.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Split out as #73859 with its own test. ✅

# 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")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

minor — this guard is only exercised by test_empty_seed_requires_zero_jitter in airflow-core/tests/unit/timetables/, so breeze testing task-sdk-tests — which runs the Task SDK distribution standalone — does not cover it. task-sdk/tests/task_sdk/definitions/timetables/test__cron.py already exists and is the natural home for it.

Copy link
Copy Markdown
Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Added to test__cron.py: empty seed, negative window, and defaults being a no-op. ✅

@potiuk

potiuk commented Sep 22, 2026

Copy link
Copy Markdown
Member

cc: @uranusjr @kaxil -> I think that one should be reviewed by you :)

@Pebble32 Pebble32 mentioned this pull request Sep 28, 2026
1 task done
@Pebble32
Pebble32 force-pushed the add-cron-timetable-jitter branch from 49838cb to f66ad85 Compare September 28, 2026 17:26
@Pebble32

Copy link
Copy Markdown
Author

Thanks for the thorough review @potiuk, all points addressed and rebased on main. description now shows the jitter offset (summary stays the bare cron), and the hash fix is split out as #73859.

On seed defaulting to the dag_id: agree the copied Dag footgun is real, I would do it as a follow up to keep this PR scoped. PTAL 🙏

@Pebble32
Pebble32 requested a review from potiuk September 28, 2026 18:00
@ashb

ashb commented Sep 29, 2026

Copy link
Copy Markdown
Member

One option thing to consider here is the Jenkins approach of using H: So H * * * * would be once per hour, as a fixed time per dag, but not the top of the hour.

cron('H 2 * * *')   // Daily around 2 AM (H spreads load)

https://crongenerator.dev/cron-jenkins/

@Pebble32

Copy link
Copy Markdown
Author

Good pointer. croniter already supports H via hash_id, so this would be small to add and would replace the offset machinery entirely. The catch is the same as before: H needs a per Dag hash_id, so the dag_id or an explicit seed.

Happy to switch this PR to H syntax if that is preferred over seed/max_jitter, or to add it as a follow up. Which would you rather?

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>
…iption, Dag naming

Signed-off-by: Adam <111773160+Pebble32@users.noreply.github.com>
@Pebble32
Pebble32 force-pushed the add-cron-timetable-jitter branch from f66ad85 to 65746f4 Compare September 30, 2026 08:12

@potiuk potiuk left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the thorough follow-up. I re-checked every point from the last round against the current diff and all of them are addressed. The partition key is back on the cron boundary, negative windows are rejected in both layers, the md5 only runs when jitter is on, the Task SDK guard is tested in test__cron.py, and description now shows the offset. #73859 has merged and is already in this branch's base, so nothing is waiting on it. Deferring "seed defaults to dag_id" to a follow-up is fine with me.

One issue remains, and it affects the main use case: enabling jitter on a Dag that already has runs.

Turning jitter on (or increasing the offset) schedules a second run for the tick that just ran. _get_next/_align_to_prev work as apply(cron(strip(t))), so a previous run_after that sits on the plain cron boundary gets stripped to before that boundary and maps back onto the same tick. I reproduced it with @daily, seed="my_dag", max_jitter=1h (offset 0:58:51) and an un-jittered last run:

CronTriggerTimetable       last run 10-01 00:00  -> next 10-01 00:58:51   (jitter off: 10-02 00:00) - second run for 10-01
CronDataIntervalTimetable  last [09-30, 10-01]   -> next [09-30 00:58:51, 10-01 00:58:51]  - overlaps the previous run by 23h

The same happens with catchup=False when jitter is switched on before that day's jittered time (last run today 00:00, now 00:25 -> next run today 00:58), and when max_jitter grows or seed changes. By the same stepping, CronPartitionTimetable gets a duplicate partition key. The current tests start either from no previous run or from a run that was already jittered, so they don't hit it. Please map a previous run_after/end to the tick it belongs to before stepping, and add a test that starts from an un-jittered last run for each of the three timetable kinds.

Smaller things:

  • trigger.py _get_partition_date: the _get_next → _get_next_cron change for run_offset != 0 isn't covered. Reverting it keeps the suite green. A _get_partition_info assertion with run_offset=1 and run_offset=-1 would cover it. The run_offset parametrization on test_iter_partition_dagrun_infos_unaffected_by_jitter doesn't exercise anything, because iter_partition_dagrun_infos never reads run_offset.
  • MultipleCronTriggerTimetable.description now repeats ", jittered by …" once per expression. Appending it once at the parent level would read better.
  • Nit: __eq__/__hash__ compare seed even when max_jitter is zero, but seed isn't serialized in that case, so such a timetable is not equal to its own round-tripped copy.

On the Jenkins-style H syntax: let's keep this PR on seed/max_jitter and take H (via croniter's hash_id) as a separate follow-up, so this one can land once the issue above is fixed.


Drafted-by: Claude Code (Opus 5); reviewed by @potiuk before posting

@Pebble32
Pebble32 requested a review from potiuk October 9, 2026 08:35

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants