Skip to content

Commit 00d736b

Browse files
committed
feat(experiments): replace the 24 schedules with one hourly schedule
1 parent e83862b commit 00d736b

4 files changed

Lines changed: 85 additions & 59 deletions

File tree

‎posthog/temporal/experiments/README.md‎

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -92,13 +92,19 @@ The `ExperimentSavedMetricsWorkflow` follows the same structure but:
9292

9393
The timeseries sync row is assembled from that day's points, so its window is approximate and a
9494
metric the daily run could not compute has no row at all. A separate workflow therefore starts a
95-
real recalculation for each eligible experiment, on its own 24 schedules
95+
real recalculation for each eligible experiment, on one hourly schedule
9696
(`products/experiments/backend/temporal/schedule.py`).
9797

98-
Those schedules fire at `:30`, after the timeseries schedules at `:00`. The coordinator selects
99-
experiments with the same rules the daily discovery uses, plus a 12-hour minimum age, an
100-
organization feature flag, and a 50-exposure floor. It then starts an ordinary
101-
`ExperimentMetricsRecalculationWorkflow` per experiment, through the same function the API uses.
98+
That schedule fires at `:30`, after the timeseries schedules at `:00`, and carries no input.
99+
Discovery reads the hour and selects the teams configured for it, so one schedule serves all 24
100+
hours. The coordinator selects experiments with the same rules the daily discovery uses, plus a
101+
12-hour minimum age, an organization feature flag, and a 50-exposure floor. It then starts an
102+
ordinary `ExperimentMetricsRecalculationWorkflow` per experiment, through the same function the
103+
API uses.
104+
105+
The timeseries workflows keep one schedule per hour, because each runs its own metric queries and
106+
can outlive its hour. This one starts other workflows and returns, so a single schedule with a
107+
`SKIP` overlap policy covers it.
102108

103109
It skips an experiment whose recalculation is already running, or whose last one finished within
104110
the hour. That freshness check ignores `timeseries_sync` rows: one is published 30 minutes before

‎products/experiments/backend/temporal/models.py‎

Lines changed: 0 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -219,13 +219,6 @@ class ExperimentPrecomputeEnrollmentCensusInputs:
219219
SCHEDULED_RECALCULATION_WORKFLOW_NAME = "experiment-scheduled-recalculation-workflow"
220220

221221

222-
@dataclasses.dataclass(frozen=False)
223-
class ScheduledRecalculationWorkflowInputs:
224-
"""Input to the scheduled recalculation coordinator."""
225-
226-
hour: int # 0-23, which hour's teams to process
227-
228-
229222
@frozen
230223
class ScheduledRecalculationStartResult:
231224
"""Outcome of one experiment's start attempt, for the coordinator's summary counts."""

‎products/experiments/backend/temporal/schedule.py‎

Lines changed: 42 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,6 @@
1919
SCHEDULED_RECALCULATION_WORKFLOW_NAME,
2020
ExperimentPrecomputeCanaryInputs,
2121
ExperimentPrecomputeEnrollmentCensusInputs,
22-
ScheduledRecalculationWorkflowInputs,
2322
)
2423

2524
CANARY_SCHEDULE_ID = "experiment-precompute-canary-schedule"
@@ -79,42 +78,57 @@ async def create_experiment_precompute_enrollment_census_schedule(client: Client
7978
await a_create_schedule(client, ENROLLMENT_CENSUS_SCHEDULE_ID, schedule, trigger_immediately=False)
8079

8180

82-
SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX = "experiment-scheduled-recalculation-hour"
81+
SCHEDULED_RECALCULATION_SCHEDULE_ID = "experiment-scheduled-recalculation"
82+
83+
# Each hour had its own schedule before the workflow resolved the hour itself. Temporal keeps a
84+
# schedule until it is deleted, so the replaced ones are removed when the single one is created.
85+
LEGACY_HOURLY_SCHEDULE_ID_PREFIX = "experiment-scheduled-recalculation-hour"
8386

8487

8588
async def create_experiment_scheduled_recalculation_schedules(client: Client) -> None:
86-
"""Create or update 24 schedules, one per hour, that start a real metrics recalculation for
87-
every eligible experiment belonging to teams configured for that hour.
89+
"""Create or update the hourly schedule that starts a real metrics recalculation for every
90+
experiment eligible in the hour it fires.
91+
92+
Fires at :30, after the daily timeseries schedule at :00, so the timeseries sync publish
93+
usually lands before a real run supersedes it. The workflow takes no input: discovery reads
94+
the hour and selects the teams configured for it, so one schedule serves all 24 hours.
8895
89-
Each fires at :30, after the daily timeseries schedule at :00, so the timeseries sync publish
90-
usually lands before a real run supersedes it. 24 separate schedules rather than one hourly
91-
schedule, matching the timeseries workflows: a run longer than an hour cannot overlap itself,
92-
so no overlap policy is needed.
96+
SKIP overlap: a run that outlives its hour must not stack on the next one.
9397
"""
94-
for hour in range(24):
95-
schedule_id = f"{SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX}-{hour:02d}"
98+
schedule = Schedule(
99+
action=ScheduleActionStartWorkflow(
100+
SCHEDULED_RECALCULATION_WORKFLOW_NAME,
101+
id=f'{SCHEDULED_RECALCULATION_SCHEDULE_ID}-{{{{.ScheduledTime.Format "2006-01-02T15"}}}}',
102+
task_queue=settings.GENERAL_PURPOSE_TASK_QUEUE,
103+
),
104+
spec=ScheduleSpec(cron_expressions=["30 * * * *"]),
105+
policy=SchedulePolicy(overlap=ScheduleOverlapPolicy.SKIP),
106+
)
107+
108+
if await a_schedule_exists(client, SCHEDULED_RECALCULATION_SCHEDULE_ID):
109+
await a_update_schedule(client, SCHEDULED_RECALCULATION_SCHEDULE_ID, schedule)
110+
else:
111+
await a_create_schedule(client, SCHEDULED_RECALCULATION_SCHEDULE_ID, schedule, trigger_immediately=False)
96112

97-
schedule = Schedule(
98-
action=ScheduleActionStartWorkflow(
99-
SCHEDULED_RECALCULATION_WORKFLOW_NAME,
100-
ScheduledRecalculationWorkflowInputs(hour=hour),
101-
id=f'{SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX}-{hour:02d}-{{{{.ScheduledTime.Format "2006-01-02"}}}}',
102-
task_queue=settings.GENERAL_PURPOSE_TASK_QUEUE,
103-
),
104-
spec=ScheduleSpec(cron_expressions=[f"30 {hour} * * *"]),
105-
policy=SchedulePolicy(overlap=ScheduleOverlapPolicy.SKIP),
106-
)
113+
await _delete_legacy_hourly_schedules(client)
107114

108-
if await a_schedule_exists(client, schedule_id):
109-
await a_update_schedule(client, schedule_id, schedule)
110-
else:
111-
await a_create_schedule(client, schedule_id, schedule, trigger_immediately=False)
112115

116+
async def _delete_legacy_hourly_schedules(client: Client) -> None:
117+
"""Remove the 24 per-hour schedules the single schedule replaces.
113118
114-
async def delete_experiment_scheduled_recalculation_schedules(client: Client) -> None:
115-
"""Delete all 24 scheduled recalculation schedules."""
119+
Without this they keep firing alongside it, starting the workflow 25 times an hour.
120+
"""
116121
for hour in range(24):
117122
try:
118-
await a_delete_schedule(client, f"{SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX}-{hour:02d}")
123+
await a_delete_schedule(client, f"{LEGACY_HOURLY_SCHEDULE_ID_PREFIX}-{hour:02d}")
119124
except Exception:
120-
pass # Schedule might not exist
125+
pass # Already gone, which is the steady state after the first run.
126+
127+
128+
async def delete_experiment_scheduled_recalculation_schedules(client: Client) -> None:
129+
"""Delete the scheduled recalculation schedule, and any legacy per-hour ones left behind."""
130+
try:
131+
await a_delete_schedule(client, SCHEDULED_RECALCULATION_SCHEDULE_ID)
132+
except Exception:
133+
pass # Schedule might not exist
134+
await _delete_legacy_hourly_schedules(client)

‎products/experiments/backend/temporal/test_scheduled_recalculation_schedule.py‎

Lines changed: 32 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,8 @@
55

66
from products.experiments.backend.temporal.models import SCHEDULED_RECALCULATION_WORKFLOW_NAME
77
from products.experiments.backend.temporal.schedule import (
8-
SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX,
8+
LEGACY_HOURLY_SCHEDULE_ID_PREFIX,
9+
SCHEDULED_RECALCULATION_SCHEDULE_ID,
910
create_experiment_scheduled_recalculation_schedules,
1011
)
1112

@@ -14,41 +15,53 @@
1415
MODULE = "products.experiments.backend.temporal.schedule"
1516

1617

17-
async def test_creates_one_schedule_per_hour_offset_from_the_timeseries_run():
18+
async def test_creates_one_hourly_schedule_offset_from_the_timeseries_run():
1819
with (
1920
patch(f"{MODULE}.a_schedule_exists", AsyncMock(return_value=False)),
2021
patch(f"{MODULE}.a_create_schedule", AsyncMock()) as create,
2122
patch(f"{MODULE}.a_update_schedule", AsyncMock()) as update,
23+
patch(f"{MODULE}.a_delete_schedule", AsyncMock()),
2224
):
2325
await create_experiment_scheduled_recalculation_schedules(AsyncMock())
2426

25-
assert create.await_count == 24
27+
assert create.await_count == 1
2628
update.assert_not_awaited()
2729

28-
schedule_ids = [call.args[1] for call in create.await_args_list]
29-
assert schedule_ids[0] == f"{SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX}-00"
30-
assert schedule_ids[23] == f"{SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX}-23"
30+
# Every field a typo could break: a wrong workflow name or task queue fails only in
31+
# production, where nothing reports it.
32+
call = create.await_args_list[0]
33+
schedule_id, schedule = call.args[1], call.args[2]
34+
assert schedule_id == SCHEDULED_RECALCULATION_SCHEDULE_ID
35+
# :30, so the daily timeseries run at :00 publishes before a real run supersedes it.
36+
assert schedule.spec.cron_expressions == ["30 * * * *"]
37+
assert schedule.action.workflow == SCHEDULED_RECALCULATION_WORKFLOW_NAME
38+
assert schedule.action.task_queue == settings.GENERAL_PURPOSE_TASK_QUEUE
39+
# No input: discovery resolves the hour, so one schedule serves every hour.
40+
assert schedule.action.args == []
3141

32-
# Every field a typo could break, checked per schedule rather than at the ends: a wrong
33-
# workflow name or task queue fails only in production, where nothing reports it.
34-
for hour, call in enumerate(create.await_args_list):
35-
schedule_id, schedule = call.args[1], call.args[2]
36-
assert schedule_id == f"{SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX}-{hour:02d}"
37-
# :30, so the daily timeseries run at :00 publishes before a real run supersedes it.
38-
assert schedule.spec.cron_expressions == [f"30 {hour} * * *"]
39-
assert schedule.action.workflow == SCHEDULED_RECALCULATION_WORKFLOW_NAME
40-
assert schedule.action.task_queue == settings.GENERAL_PURPOSE_TASK_QUEUE
41-
assert schedule.action.args[0].hour == hour
42-
assert schedule.action.id.startswith(f"{SCHEDULED_RECALCULATION_SCHEDULE_ID_PREFIX}-{hour:02d}-")
42+
43+
async def test_the_replaced_per_hour_schedules_are_deleted():
44+
# Temporal keeps a schedule until it is deleted, so leaving the 24 behind would start the
45+
# workflow 25 times an hour.
46+
with (
47+
patch(f"{MODULE}.a_schedule_exists", AsyncMock(return_value=False)),
48+
patch(f"{MODULE}.a_create_schedule", AsyncMock()),
49+
patch(f"{MODULE}.a_delete_schedule", AsyncMock()) as delete,
50+
):
51+
await create_experiment_scheduled_recalculation_schedules(AsyncMock())
52+
53+
deleted = [call.args[1] for call in delete.await_args_list]
54+
assert deleted == [f"{LEGACY_HOURLY_SCHEDULE_ID_PREFIX}-{hour:02d}" for hour in range(24)]
4355

4456

45-
async def test_existing_schedules_are_updated_not_recreated():
57+
async def test_an_existing_schedule_is_updated_not_recreated():
4658
with (
4759
patch(f"{MODULE}.a_schedule_exists", AsyncMock(return_value=True)),
4860
patch(f"{MODULE}.a_create_schedule", AsyncMock()) as create,
4961
patch(f"{MODULE}.a_update_schedule", AsyncMock()) as update,
62+
patch(f"{MODULE}.a_delete_schedule", AsyncMock()),
5063
):
5164
await create_experiment_scheduled_recalculation_schedules(AsyncMock())
5265

53-
assert update.await_count == 24
66+
assert update.await_count == 1
5467
create.assert_not_awaited()

0 commit comments

Comments
 (0)