Skip to content

Commit d06f11d

Browse files
breken-aiclaude
andcommitted
fix(scheduler): skip intervals that overlap unready parent intervals
`interval_diff` advanced past an unready interval whenever the candidate interval started before it. When a daily model depends on an hourly model whose signal marks an hour as not ready, the day containing that hour was compared against the next unready interval instead, so the daily model was evaluated over incomplete upstream data and the interval was recorded as done. Ready intervals after the last overlap were also dropped. Always skip the candidate interval on overlap and keep every remaining interval once the unready intervals are exhausted. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Signed-off-by: breken-ai <312387581+breken-ai@users.noreply.github.com>
1 parent 263723f commit d06f11d

2 files changed

Lines changed: 83 additions & 8 deletions

File tree

‎sqlmesh/core/scheduler.py‎

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1158,15 +1158,12 @@ def interval_diff(
11581158
if uninterrupted:
11591159
return results
11601160

1161-
if interval_a[0] >= interval_b[0]:
1162-
index_a += 1
1163-
else:
1164-
index_b += 1
1161+
# interval_a overlaps interval_b, so it must be excluded. interval_b is kept because
1162+
# it may also overlap the next interval in intervals_a.
1163+
index_a += 1
11651164

1166-
if index_a < len_a:
1167-
interval_a = intervals_a[index_a]
1168-
if interval_a[0] >= interval_b[1] or interval_b[0] >= interval_a[1]:
1169-
results.extend(intervals_a[index_a:])
1165+
# Any remaining intervals in intervals_a start after the last interval in intervals_b
1166+
results.extend(intervals_a[index_a:])
11701167

11711168
return results
11721169

‎tests/core/test_scheduler.py‎

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -677,6 +677,15 @@ def test_interval_diff():
677677
assert interval_diff([(1, 2), (2, 3)], [(3, 4)], uninterrupted=True) == [(1, 2), (2, 3)]
678678
assert interval_diff([(1, 2), (2, 3)], [(2, 3)], uninterrupted=True) == [(1, 2)]
679679

680+
# An interval that contains a smaller interval from `intervals_b` (for example a daily
681+
# interval that contains an unready hourly parent interval) must be excluded, and the
682+
# intervals after it must still be kept.
683+
assert interval_diff([(0, 24), (24, 48)], [(10, 11)]) == [(24, 48)]
684+
assert interval_diff(
685+
[(0, 24), (24, 48), (48, 72), (72, 96)],
686+
[(10, 11), (53, 54)],
687+
) == [(24, 48), (72, 96)]
688+
680689

681690
def test_signal_intervals(mocker: MockerFixture, make_snapshot, get_batched_missing_intervals):
682691
@signal()
@@ -870,6 +879,75 @@ def signal_base(batch: DatetimeRanges):
870879
}
871880

872881

882+
def test_signals_unready_hourly_parent_blocks_daily_child(
883+
mocker: MockerFixture, make_snapshot, get_batched_missing_intervals
884+
):
885+
@signal()
886+
def hourly_ready(batch: DatetimeRanges):
887+
# 2023-01-01 10:00 and 2023-01-03 05:00 are not ready yet
888+
return [i for i in batch if (i[0].day, i[0].hour) not in ((1, 10), (3, 5))]
889+
890+
signals = signal.get_registry()
891+
892+
parent = make_snapshot(
893+
load_sql_based_model(
894+
parse( # type: ignore
895+
"""
896+
MODEL (
897+
name hourly_parent,
898+
kind INCREMENTAL_BY_TIME_RANGE(
899+
time_column dt,
900+
),
901+
cron '@hourly',
902+
start '2023-01-01',
903+
signals HOURLY_READY(),
904+
);
905+
SELECT @start_dt AS dt;
906+
"""
907+
),
908+
signal_definitions=signals,
909+
),
910+
)
911+
912+
child = make_snapshot(
913+
load_sql_based_model(
914+
parse( # type: ignore
915+
"""
916+
MODEL (
917+
name daily_child,
918+
kind INCREMENTAL_BY_TIME_RANGE(
919+
time_column dt,
920+
),
921+
start '2023-01-01',
922+
);
923+
SELECT dt FROM hourly_parent
924+
"""
925+
),
926+
signal_definitions=signals,
927+
),
928+
nodes={parent.name: parent.model},
929+
)
930+
931+
snapshot_evaluator = SnapshotEvaluator(adapters=mocker.MagicMock(), ddl_concurrent_tasks=1)
932+
scheduler = Scheduler(
933+
snapshots=[parent, child],
934+
snapshot_evaluator=snapshot_evaluator,
935+
state_sync=mocker.MagicMock(),
936+
max_workers=2,
937+
default_catalog=None,
938+
)
939+
940+
batches = get_batched_missing_intervals(
941+
scheduler, "2023-01-01", "2023-01-04", "2023-01-05 01:00:00"
942+
)
943+
944+
# Days containing an unready parent hour must not be evaluated, the other days must be
945+
assert batches[child] == [
946+
(to_timestamp("2023-01-02"), to_timestamp("2023-01-03")),
947+
(to_timestamp("2023-01-04"), to_timestamp("2023-01-05")),
948+
]
949+
950+
873951
@pytest.mark.parametrize(
874952
"batch_size, expected_batches",
875953
[

0 commit comments

Comments
 (0)