Repository navigation
Conversation
9121d44 to
49838cb
Compare
potiuk
left a comment
There was a problem hiding this comment.
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/descriptiondon'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 whethersummaryshould surface the offset — otherwise the schedule shown and the schedule run disagree, which is a support question waiting to happen.seedis 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 defaultseedto 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 ownstimetables/@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
| 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") |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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) |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Moved inside the jitter branch. ✅
| """ | ||
| Mixin to provide interface to work with croniter. | ||
|
|
||
| Optionally applies a deterministic, per-DAG jitter to every scheduled time. |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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)) |
There was a problem hiding this comment.
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.
| # 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") |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
Added to test__cron.py: empty seed, negative window, and defaults being a no-op. ✅
49838cb to
f66ad85
Compare
|
Thanks for the thorough review @potiuk, all points addressed and rebased on main. 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 🙏 |
|
One option thing to consider here is the Jenkins approach of using cron('H 2 * * *') // Daily around 2 AM (H spreads load) |
|
Good pointer. croniter already supports Happy to switch this PR to |
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>
f66ad85 to
65746f4
Compare
potiuk
left a comment
There was a problem hiding this comment.
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_cronchange forrun_offset != 0isn't covered. Reverting it keeps the suite green. A_get_partition_infoassertion withrun_offset=1andrun_offset=-1would cover it. Therun_offsetparametrization ontest_iter_partition_dagrun_infos_unaffected_by_jitterdoesn't exercise anything, becauseiter_partition_dagrun_infosnever readsrun_offset.MultipleCronTriggerTimetable.descriptionnow repeats ", jittered by …" once per expression. Appending it once at the parent level would read better.- Nit:
__eq__/__hash__compareseedeven whenmax_jitteris zero, butseedisn'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
Add opt-in, deterministic jitter to every cron-based timetable via
CronMixin: two kw-only params,seedandmax_jitter, onCronTriggerTimetable,CronDataIntervalTimetable,MultipleCronTriggerTimetableandCronPartitionTimetable, 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
CronMixinso all cron scheduling inherits it. It supersedes #69705, which I will close once this lands. cc @kaxil, who reviewed the original.Why
@dailyexpands to0 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:@dailyintentparallelism,max_active_tasks_per_dag, pool slots)Jitter is peak shaving on top of concurrency control, not a replacement for it. It is the same idea as Jenkins'
Hcron syntax. The offset is deterministic: the sameseed(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_jitteradded toCronMixinin both layers, so all four cron timetables inherit them; the concrete classes thread the params through their constructors andserialize/deserialize, and theencoders.pyvariants emit them.md5(seed) % max_jitter, computed in integer microseconds (FIPS-safehashlib_wrapper.md5), and applied as a "strip → cron → apply" coordinate shift inCronMixin._get_next/_get_prev. Because_align_to_next/_align_to_prevbuild on those primitives, every cron timetable inherits the shift with no per-class scheduling code.max_jitterat 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
max_jittersmall relative to the period.MultipleCronTriggerTimetablechildren share one offset. The sameseed/max_jitteris passed to every child, so the whole DAG shifts in lockstep and its fire times keep their relative spacing.seed/max_jitterwhen jitter is set. The wire format of existing DAGs is unchanged (so nothing is re-serialized on upgrade) anddeserializetolerates missing keys. The existing serialization tests pass untouched.__eq__/__hash__, since timetables differing only in jitter produce different schedules. While adding this I fixed a latent bug:__hash__included the pendulumTimezoneobject, which is unhashable, sohash()on any cron timetable raised; it now hashesstr(timezone).seedwithmax_jitter > 0raises 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 existingtest_trigger_timetable.py:America/New_York);max_jitterreproduces the plain timetable exactly (catchup on/off);[0, max_jitter), spread across distinct seeds, and correct for sub-second windows;serialize/deserializeand SDK encode → decode round-trips preserve the offset for all four timetables;MultipleCronTriggerTimetablechildren share a single offset;Was generative AI tooling used to co-author this PR?
Generated-by: Claude (Claude Code), following the guidelines
related: #69705
related: #69027