From fb7dc873530d95379da7f553ba3f997df8600534 Mon Sep 17 00:00:00 2001 From: Kyle Consalus Date: Mon, 21 Sep 2026 12:00:22 -0700 Subject: [PATCH 1/5] WIP: Probable --- src/sentry/issues/derived/tasks.py | 42 ++++++- src/sentry/issues/derived/tasks_util.py | 132 +++++++++++----------- tests/sentry/issues/derived/test_tasks.py | 117 +++++++++++++++---- 3 files changed, 204 insertions(+), 87 deletions(-) diff --git a/src/sentry/issues/derived/tasks.py b/src/sentry/issues/derived/tasks.py index 5e4a4fbacebe..890b5f0dd66e 100644 --- a/src/sentry/issues/derived/tasks.py +++ b/src/sentry/issues/derived/tasks.py @@ -581,7 +581,11 @@ def heal_stale_derived_data(**kwargs: object) -> None: lower_bound = 0 if stale_hash is None else state.stale[stale_hash] logger.info( "heal_stale_derived_data.range_selection_started", - extra={"hash_kind": hash_kind, "remaining_budget": remaining}, + extra={ + "hash_kind": hash_kind, + "remaining_budget": remaining, + "group_id_lower_bound": lower_bound, + }, ) range_selection_started_at = time.monotonic() try: @@ -597,6 +601,8 @@ def heal_stale_derived_data(**kwargs: object) -> None: extra={ "hash_kind": hash_kind, "elapsed": time.monotonic() - range_selection_started_at, + "pipeline_hash": stale_hash, + "group_id_lower_bound": lower_bound, }, ) metrics.incr( @@ -911,6 +917,7 @@ def regenerate_stale_derived_data_batch( ) from taskbroker_client.state import current_task + from sentry import options from sentry.issues.derived.promote import build_and_promote_batch from sentry.issues.models.groupderiveddata import GroupDerivedData from sentry.taskworker.selfchain_idempotency import already_spawned, mark_spawned @@ -944,6 +951,7 @@ def regenerate_stale_derived_data_batch( start = time.monotonic() + batch_size = max(1, options.get("issues.derived.heal-batch-size")) group_ids = list( GroupDerivedData.objects.filter( pipeline_hash=target_hash, # a None target renders as IS NULL @@ -951,7 +959,22 @@ def regenerate_stale_derived_data_batch( group_id__lt=group_id_end, ) .order_by("group_id") - .values_list("group_id", flat=True) + .values_list("group_id", flat=True)[: batch_size + 1] + ) + range_overflow = len(group_ids) > batch_size + if range_overflow: + group_ids = group_ids[:batch_size] + + # How close the scheduler's density estimate landed to reality. Censored at + # batch_size, so read it together with the ``range_overflow`` reschedule + # counter: a distribution well under batch_size means ranges are too wide and + # slots are being wasted, while frequent overflow means they are too narrow. + # These two are what the density sampling constants should be tuned from. + metrics.distribution( + "issues.derived.heal_range_rows_found", + len(group_ids), + sample_rate=1.0, + tags={"hash_kind": "null" if target_hash is None else "stale"}, ) result = build_and_promote_batch( @@ -981,6 +1004,21 @@ def regenerate_stale_derived_data_batch( ) if activation_id: mark_spawned(_REGENERATE_STALE_BATCH_TASK_KEY, activation_id) + elif range_overflow: + rescheduled = True + metrics.incr( + "issues.derived.regenerate_stale_batch_rescheduled", + sample_rate=1.0, + tags={"reason": "range_overflow"}, + ) + regenerate_stale_derived_data_batch.delay( + stale_pipeline_hashes=stale_pipeline_hashes, + target_hash=target_hash, + group_id_start=group_ids[-1] + 1, + group_id_end=group_id_end, + ) + if activation_id: + mark_spawned(_REGENERATE_STALE_BATCH_TASK_KEY, activation_id) _record_batch_metrics( result.processed, diff --git a/src/sentry/issues/derived/tasks_util.py b/src/sentry/issues/derived/tasks_util.py index 5da712e5c530..70def955c8d5 100644 --- a/src/sentry/issues/derived/tasks_util.py +++ b/src/sentry/issues/derived/tasks_util.py @@ -4,6 +4,7 @@ import random from dataclasses import dataclass from datetime import datetime, timedelta, timezone +from math import ceil from typing import Protocol from django.db import connections, router @@ -18,11 +19,16 @@ logger = logging.getLogger(__name__) _MAX_CHECK_GROUPS = 10_000 - -# Safety valve on the number of group IDs one ``group_id_ranges_for_hash`` call may -# walk, however large the requested chunking is. -_MAX_SCANNED_GROUP_IDS = 2_000_000 -_GROUP_ID_RANGE_TIMEOUT = timedelta(seconds=50) +_GROUP_ID_RANGE_QUERY_TIMEOUT = timedelta(seconds=40) +# Rows read per density probe, and the number of probes per call. Their product +# bounds the scan; the ratio of scheduled groups to rows read is also exactly how +# far each probe is extrapolated, so these trade query cost against estimate error. +# Sample size fights random gap variance (error falls as 1/sqrt(size), so there is +# little gained past ~100); probe count fights systematic density drift across the +# range, which it shortens linearly. Tune from +# ``issues.derived.heal_range_rows_found`` rather than from first principles. +_RANGE_DENSITY_SAMPLE_SIZE = 100 +_MAX_RANGE_DENSITY_SAMPLES = 40 @dataclass(frozen=True) @@ -117,81 +123,75 @@ def _pick_random_fresh_group_ranges( def group_id_ranges_for_hash( pipeline_hash: str | None, *, chunk_size: int, max_chunks: int, group_id_lower_bound: int = 0 ) -> GroupIdRangeResult: - """Partition the group IDs of GroupDerivedData rows with a pipeline_hash into ranges. + """Estimate ranges covering GroupDerivedData rows with a pipeline_hash. Returns at most max_chunks of ascending disjoint [start, end) ranges, each - covering chunk_size group IDs but possibly the last. ``drained`` is true only - when a valid query found no rows at or above ``group_id_lower_bound``. + targeting chunk_size group IDs. Density is sampled at a bounded number of + points so the query cost stays flat as the scheduling budget grows, which + makes the per-region throughput knob independent of how long this runs. + Ranges are therefore approximate: an over-dense one is split by the worker, + and an under-dense one costs only a scheduling slot. ``drained`` is true + only when a valid query found no rows at or above ``group_id_lower_bound``. + + Raises ``OperationalError`` if all density probes together exceed the server-side + statement timeout. Callers must not interpret that as a drained hash. """ if chunk_size <= 0 or max_chunks <= 0: return GroupIdRangeResult(ranges=[], drained=False) - # One boundary per chunk, plus one extra to close the final range (or, if we ran - # out of matching rows first, to tell us we did). - scan_limit = chunk_size * (max_chunks + 1) - if scan_limit > _MAX_SCANNED_GROUP_IDS: - logger.warning( - "group_id_ranges_for_hash.scan_budget_clamped", - extra={ - "chunk_size": chunk_size, - "max_chunks": max_chunks, - "requested_scan_limit": scan_limit, - "max_scan_limit": _MAX_SCANNED_GROUP_IDS, - }, - ) - # Never clamp below two chunks: a lone boundary can't close a range, so we'd - # return nothing and the caller would take that to mean there was nothing to do. - scan_limit = max(_MAX_SCANNED_GROUP_IDS, chunk_size * 2) - hash_predicate = "pipeline_hash IS NULL" if pipeline_hash is None else "pipeline_hash = %s" - # The LIMIT lives in the innermost subquery to encourage Postgres to enforce - # the limit before doing the window function stuff. ``cnt`` tells us whether we - # hit the limit, and pulls out the last row we saw so we can close the final - # chunk without a second query. sql = f""" - SELECT group_id, rn, cnt FROM ( - SELECT group_id, - -- Ordered so the numbering is defined rather than dependent on the - -- order the subquery happens to emit. count(*) stays unordered; an - -- ORDER BY there would turn it into a running count. - row_number() OVER (ORDER BY group_id) - 1 AS rn, - count(*) OVER () AS cnt - FROM ( - SELECT group_id - FROM {GroupDerivedData._meta.db_table} - WHERE {hash_predicate} AND group_id >= %s - ORDER BY group_id - LIMIT %s - ) scanned - ) numbered - WHERE mod(rn, %s) = 0 OR rn = cnt - 1 - ORDER BY rn + SELECT group_id + FROM {GroupDerivedData._meta.db_table} + WHERE {hash_predicate} AND group_id >= %s + ORDER BY pipeline_hash, group_id + LIMIT %s """ - params: list[str | int] = [] if pipeline_hash is None else [pipeline_hash] - params += [group_id_lower_bound, scan_limit, chunk_size] using = router.db_for_read(GroupDerivedData) + ranges: list[tuple[int, int]] = [] + next_group_id = group_id_lower_bound + sample_limit = _RANGE_DENSITY_SAMPLE_SIZE + 1 + sample_count = min(max_chunks, _MAX_RANGE_DENSITY_SAMPLES) + with ( - statement_timeout(using, _GROUP_ID_RANGE_TIMEOUT), metrics.timer("issues.derived.group_id_range_query"), - connections[using].cursor() as cursor, + statement_timeout(using, _GROUP_ID_RANGE_QUERY_TIMEOUT), + connections[using].cursor() as db_cursor, ): - cursor.execute(sql, params) - rows = cursor.fetchall() - - if not rows: - return GroupIdRangeResult(ranges=[], drained=True) - - scanned = rows[0][2] - boundaries = [group_id for group_id, rn, _ in rows if rn % chunk_size == 0] - - if scanned == scan_limit: - # We got more than enough rows; the trailing boundary closes the last range. - ends = boundaries[1:] - else: - # We didn't fill out the final chunk, so the last row we scanned closes it. - ends = boundaries[1:] + [rows[-1][0] + 1] - return GroupIdRangeResult(ranges=list(zip(boundaries, ends))[:max_chunks], drained=False) + while len(ranges) < max_chunks: + params: list[str | int] = [] if pipeline_hash is None else [pipeline_hash] + params += [next_group_id, sample_limit] + db_cursor.execute(sql, params) + sampled_group_ids = [row[0] for row in db_cursor.fetchall()] + + if not sampled_group_ids: + return GroupIdRangeResult(ranges=ranges, drained=not ranges) + + if len(sampled_group_ids) <= _RANGE_DENSITY_SAMPLE_SIZE: + starts = sampled_group_ids[::chunk_size] + ends = starts[1:] + [sampled_group_ids[-1] + 1] + ranges.extend(zip(starts, ends)) + break + + density_sample = sampled_group_ids[:_RANGE_DENSITY_SAMPLE_SIZE] + sample_width = density_sample[-1] - density_sample[0] + 1 + estimated_chunk_width = max( + 1, + ceil(sample_width * chunk_size / len(density_sample)), + ) + samples_left = min(sample_count, max_chunks - len(ranges)) + chunks_for_sample = ceil((max_chunks - len(ranges)) / samples_left) + + range_start = density_sample[0] + for _ in range(chunks_for_sample): + range_end = range_start + estimated_chunk_width + ranges.append((range_start, range_end)) + range_start = range_end + next_group_id = range_start + sample_count -= 1 + + return GroupIdRangeResult(ranges=ranges[:max_chunks], drained=False) def _resume_check_id( diff --git a/tests/sentry/issues/derived/test_tasks.py b/tests/sentry/issues/derived/test_tasks.py index 5769ce118f75..338f1dcf6913 100644 --- a/tests/sentry/issues/derived/test_tasks.py +++ b/tests/sentry/issues/derived/test_tasks.py @@ -645,6 +645,42 @@ def test_range_selection_timeout_is_reported_and_other_hashes_continue(self) -> in mock_incr.call_args_list ) + def test_range_selection_timeout_does_not_advance_mark(self) -> None: + stale = self._pick_stale_hash() + state = HealSchedulerState( + head_hash=PIPELINE.pipeline_hash, + stale={stale: 10}, + discovered_at=datetime.now(timezone.utc), + ) + with ( + override_options( + { + "issues.derived.heal-max-tasks": 1, + "issues.derived.check-task-count": 0, + } + ), + patch("sentry.issues.derived.tasks.load_state", return_value=state), + patch( + "sentry.issues.derived.tasks_util.group_id_ranges_for_hash", + side_effect=[ + GroupIdRangeResult(ranges=[], drained=True), + OperationalError, + ], + ), + patch.object(regenerate_stale_derived_data_batch, "delay") as delay, + patch("sentry.issues.derived.tasks.logger") as mock_logger, + ): + heal_stale_derived_data() + + assert state.stale == {stale: 10} + delay.assert_not_called() + failure_log = mock_logger.exception.call_args + assert failure_log.args == ("heal_stale_derived_data.range_selection_failed",) + assert failure_log.kwargs["extra"]["hash_kind"] == "stale" + assert failure_log.kwargs["extra"]["pipeline_hash"] == stale + assert failure_log.kwargs["extra"]["group_id_lower_bound"] == 10 + assert failure_log.kwargs["extra"]["elapsed"] >= 0 + def test_discovered_hashes_are_saved_before_range_selection(self) -> None: stale_hash = self._pick_stale_hash() with ( @@ -1432,9 +1468,7 @@ def test_truncates_to_max_chunks(self) -> None: (group_ids[2], group_ids[4]), ] - def test_truncates_when_scan_limit_is_reached(self) -> None: - # 6 rows exactly fills the scan budget of chunk_size * (max_chunks + 1), so - # the tail is known to be incomplete and waits for the next call. + def test_truncates_full_tail_to_max_chunks(self) -> None: group_ids = self._seed(6, self.HASH) assert group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=2).ranges == [ @@ -1452,28 +1486,37 @@ def test_invalid_chunking(self) -> None: self.HASH, chunk_size=2, max_chunks=0 ) == GroupIdRangeResult(ranges=[], drained=False) - def test_clamps_scan_budget(self) -> None: - group_ids = self._seed(3, self.HASH) + def test_samples_local_density_for_approximate_ranges(self) -> None: + groups = self.create_unprocessed_groups(81) + all_group_ids = sorted(group.id for group in groups) + matching_group_ids = [all_group_ids[index] for index in range(0, 81, 10)] + for group_id in all_group_ids: + GroupDerivedData.objects.create( + group_id=group_id, + pipeline_hash=self.HASH if group_id in matching_group_ids else self.OTHER_HASH, + ) - with patch("sentry.issues.derived.tasks_util._MAX_SCANNED_GROUP_IDS", 2): - result = group_id_ranges_for_hash(self.HASH, chunk_size=1, max_chunks=5) + with ( + patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2), + patch("sentry.issues.derived.tasks_util._MAX_RANGE_DENSITY_SAMPLES", 2), + ): + result = group_id_ranges_for_hash(self.HASH, chunk_size=4, max_chunks=4) - # Only 2 rows were scanned, so the clamp must not be mistaken for having - # reached the end — the tail may not run past what we scanned. - assert result.ranges == [(group_ids[0], group_ids[1])] + assert len(result.ranges) == 4 + assert result.ranges == sorted(result.ranges) + assert all( + any(start <= group_id < end for start, end in result.ranges) + for group_id in matching_group_ids + ) - def test_clamp_never_drops_below_one_chunk(self) -> None: - group_ids = self._seed(3, self.HASH) + def test_exact_tail_does_not_waste_remaining_slots(self) -> None: + group_ids = self._seed(5, self.HASH) - # Clamping below chunk_size would leave a single boundary, which can't close - # a range — the caller would read the empty result as "nothing to do". - with patch("sentry.issues.derived.tasks_util._MAX_SCANNED_GROUP_IDS", 1): + with patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2): result = group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) - assert result.ranges == [ - (group_ids[0], group_ids[2]), - (group_ids[2], group_ids[2] + 1), - ] + assert result.ranges[-1][1] == group_ids[-1] + 1 + assert len(result.ranges) == 3 def test_lower_bound(self) -> None: group_ids = self._seed(3, self.HASH) @@ -1642,6 +1685,42 @@ def test_reschedules_on_batch_timeout(self) -> None: assert kwargs["group_id_start"] == group_ids[0] + 1 assert kwargs["group_id_end"] == group_ids[-1] + 1 + def test_reschedules_when_estimated_range_is_overfull(self) -> None: + groups = self.create_unprocessed_groups(3) + group_ids = sorted(g.id for g in groups) + for gid in group_ids: + process_group_log(gid) + + stale = self._stale() + GroupDerivedData.objects.filter(group_id__in=group_ids).update(pipeline_hash=stale) + + with ( + override_options({"issues.derived.heal-batch-size": 2}), + patch.object(regenerate_stale_derived_data_batch, "delay") as mock_delay, + ): + regenerate_stale_derived_data_batch( + stale_pipeline_hashes=[stale], + target_hash=stale, + group_id_start=group_ids[0], + group_id_end=group_ids[-1] + 1, + ) + + assert ( + GroupDerivedData.objects.get(group_id=group_ids[0]).pipeline_hash + == PIPELINE.pipeline_hash + ) + assert ( + GroupDerivedData.objects.get(group_id=group_ids[1]).pipeline_hash + == PIPELINE.pipeline_hash + ) + assert GroupDerivedData.objects.get(group_id=group_ids[2]).pipeline_hash == stale + mock_delay.assert_called_once_with( + stale_pipeline_hashes=[stale], + target_hash=stale, + group_id_start=group_ids[1] + 1, + group_id_end=group_ids[-1] + 1, + ) + def test_reschedules_on_group_log_timeout(self) -> None: groups = self.create_unprocessed_groups(2) group_ids = sorted(g.id for g in groups) From eabaa72190d9b502230a4a38ae63074f154526b0 Mon Sep 17 00:00:00 2001 From: Kyle Consalus Date: Tue, 22 Sep 2026 13:12:56 -0700 Subject: [PATCH 2/5] more --- src/sentry/issues/derived/tasks.py | 10 +-- src/sentry/issues/derived/tasks_util.py | 84 ++++++++++++++--------- tests/sentry/issues/derived/test_tasks.py | 27 ++++++-- 3 files changed, 81 insertions(+), 40 deletions(-) diff --git a/src/sentry/issues/derived/tasks.py b/src/sentry/issues/derived/tasks.py index 890b5f0dd66e..a5ac25e9ed41 100644 --- a/src/sentry/issues/derived/tasks.py +++ b/src/sentry/issues/derived/tasks.py @@ -965,11 +965,11 @@ def regenerate_stale_derived_data_batch( if range_overflow: group_ids = group_ids[:batch_size] - # How close the scheduler's density estimate landed to reality. Censored at - # batch_size, so read it together with the ``range_overflow`` reschedule - # counter: a distribution well under batch_size means ranges are too wide and - # slots are being wasted, while frequent overflow means they are too narrow. - # These two are what the density sampling constants should be tuned from. + # How close the scheduler's density estimate landed to reality. The worker + # records at most batch_size rows here and reports denser ranges separately as + # ``range_overflow`` reschedules. Counts well below batch_size mean ranges are + # too wide and slots are being wasted; frequent overflow means they are too + # narrow. Use both signals to tune the density sampling constants. metrics.distribution( "issues.derived.heal_range_rows_found", len(group_ids), diff --git a/src/sentry/issues/derived/tasks_util.py b/src/sentry/issues/derived/tasks_util.py index 70def955c8d5..d0a96769a0c8 100644 --- a/src/sentry/issues/derived/tasks_util.py +++ b/src/sentry/issues/derived/tasks_util.py @@ -2,6 +2,7 @@ import logging import random +import time from dataclasses import dataclass from datetime import datetime, timedelta, timezone from math import ceil @@ -9,6 +10,7 @@ from django.db import connections, router from django.db.models import Max, Min +from django.db.utils import OperationalError from sentry.issues.derived.check import CheckFailure, CheckId, CheckInvalidated, CheckResult from sentry.issues.models.groupderiveddata import GroupDerivedData @@ -20,15 +22,10 @@ _MAX_CHECK_GROUPS = 10_000 _GROUP_ID_RANGE_QUERY_TIMEOUT = timedelta(seconds=40) -# Rows read per density probe, and the number of probes per call. Their product -# bounds the scan; the ratio of scheduled groups to rows read is also exactly how -# far each probe is extrapolated, so these trade query cost against estimate error. -# Sample size fights random gap variance (error falls as 1/sqrt(size), so there is -# little gained past ~100); probe count fights systematic density drift across the -# range, which it shortens linearly. Tune from -# ``issues.derived.heal_range_rows_found`` rather than from first principles. +_MAX_EXACT_RANGE_ROWS = 10_000 _RANGE_DENSITY_SAMPLE_SIZE = 100 -_MAX_RANGE_DENSITY_SAMPLES = 40 +_RANGES_PER_DENSITY_SAMPLE = 5 +_MAX_RANGE_DENSITY_SAMPLES = 200 @dataclass(frozen=True) @@ -126,15 +123,15 @@ def group_id_ranges_for_hash( """Estimate ranges covering GroupDerivedData rows with a pipeline_hash. Returns at most max_chunks of ascending disjoint [start, end) ranges, each - targeting chunk_size group IDs. Density is sampled at a bounded number of - points so the query cost stays flat as the scheduling budget grows, which - makes the per-region throughput knob independent of how long this runs. - Ranges are therefore approximate: an over-dense one is split by the worker, - and an under-dense one costs only a scheduling slot. ``drained`` is true - only when a valid query found no rows at or above ``group_id_lower_bound``. - - Raises ``OperationalError`` if all density probes together exceed the server-side - statement timeout. Callers must not interpret that as a drained hash. + targeting chunk_size group IDs. Small requests use exact boundaries. Larger + requests sample local density once per five ranges, so estimate quality stays + stable as a region increases its scheduling budget. An over-dense range is + split by the worker, and an under-dense one costs only a scheduling slot. + ``drained`` is true only when a valid query found no rows at or above + ``group_id_lower_bound``. + + Raises ``OperationalError`` if all density probes together exceed the query + budget. Callers must not interpret that as a drained hash. """ if chunk_size <= 0 or max_chunks <= 0: return GroupIdRangeResult(ranges=[], drained=False) @@ -149,21 +146,46 @@ def group_id_ranges_for_hash( """ using = router.db_for_read(GroupDerivedData) - ranges: list[tuple[int, int]] = [] - next_group_id = group_id_lower_bound - sample_limit = _RANGE_DENSITY_SAMPLE_SIZE + 1 - sample_count = min(max_chunks, _MAX_RANGE_DENSITY_SAMPLES) - - with ( - metrics.timer("issues.derived.group_id_range_query"), - statement_timeout(using, _GROUP_ID_RANGE_QUERY_TIMEOUT), - connections[using].cursor() as db_cursor, - ): - while len(ranges) < max_chunks: + query_deadline = time.monotonic() + _GROUP_ID_RANGE_QUERY_TIMEOUT.total_seconds() + + with metrics.timer("issues.derived.group_id_range_query"): + + def fetch_group_ids(start: int, limit: int) -> list[int]: + remaining_seconds = query_deadline - time.monotonic() + if remaining_seconds <= 0.001: + raise OperationalError("group ID range query budget exceeded") + params: list[str | int] = [] if pipeline_hash is None else [pipeline_hash] - params += [next_group_id, sample_limit] - db_cursor.execute(sql, params) - sampled_group_ids = [row[0] for row in db_cursor.fetchall()] + params += [start, limit] + with ( + statement_timeout(using, timedelta(seconds=remaining_seconds)), + connections[using].cursor() as db_cursor, + ): + db_cursor.execute(sql, params) + return [row[0] for row in db_cursor.fetchall()] + + requested_rows = chunk_size * max_chunks + if requested_rows <= _MAX_EXACT_RANGE_ROWS: + group_ids = fetch_group_ids(group_id_lower_bound, requested_rows + 1) + if not group_ids: + return GroupIdRangeResult(ranges=[], drained=True) + + starts = group_ids[:requested_rows:chunk_size] + if len(group_ids) > requested_rows: + ends = starts[1:] + [group_ids[requested_rows]] + else: + ends = starts[1:] + [group_ids[-1] + 1] + return GroupIdRangeResult(ranges=list(zip(starts, ends)), drained=False) + + ranges: list[tuple[int, int]] = [] + next_group_id = group_id_lower_bound + sample_limit = _RANGE_DENSITY_SAMPLE_SIZE + 1 + sample_count = min( + ceil(max_chunks / _RANGES_PER_DENSITY_SAMPLE), + _MAX_RANGE_DENSITY_SAMPLES, + ) + while len(ranges) < max_chunks: + sampled_group_ids = fetch_group_ids(next_group_id, sample_limit) if not sampled_group_ids: return GroupIdRangeResult(ranges=ranges, drained=not ranges) diff --git a/tests/sentry/issues/derived/test_tasks.py b/tests/sentry/issues/derived/test_tasks.py index 338f1dcf6913..030346a3c490 100644 --- a/tests/sentry/issues/derived/test_tasks.py +++ b/tests/sentry/issues/derived/test_tasks.py @@ -1486,6 +1486,17 @@ def test_invalid_chunking(self) -> None: self.HASH, chunk_size=2, max_chunks=0 ) == GroupIdRangeResult(ranges=[], drained=False) + def test_density_probes_share_one_query_budget(self) -> None: + with ( + patch( + "sentry.issues.derived.tasks_util.time.monotonic", + # Query deadline, metrics timer start, budget check, metrics timer end. + side_effect=[0.0, 0.0, 41.0, 41.0], + ), + pytest.raises(OperationalError, match="query budget exceeded"), + ): + group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) + def test_samples_local_density_for_approximate_ranges(self) -> None: groups = self.create_unprocessed_groups(81) all_group_ids = sorted(group.id for group in groups) @@ -1497,7 +1508,9 @@ def test_samples_local_density_for_approximate_ranges(self) -> None: ) with ( + patch("sentry.issues.derived.tasks_util._MAX_EXACT_RANGE_ROWS", 0), patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2), + patch("sentry.issues.derived.tasks_util._RANGES_PER_DENSITY_SAMPLE", 2), patch("sentry.issues.derived.tasks_util._MAX_RANGE_DENSITY_SAMPLES", 2), ): result = group_id_ranges_for_hash(self.HASH, chunk_size=4, max_chunks=4) @@ -1509,14 +1522,20 @@ def test_samples_local_density_for_approximate_ranges(self) -> None: for group_id in matching_group_ids ) - def test_exact_tail_does_not_waste_remaining_slots(self) -> None: + def test_low_volume_request_uses_exact_boundaries(self) -> None: group_ids = self._seed(5, self.HASH) - with patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2): + with ( + patch("sentry.issues.derived.tasks_util._MAX_EXACT_RANGE_ROWS", 10), + patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2), + ): result = group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) - assert result.ranges[-1][1] == group_ids[-1] + 1 - assert len(result.ranges) == 3 + assert result.ranges == [ + (group_ids[0], group_ids[2]), + (group_ids[2], group_ids[4]), + (group_ids[4], group_ids[4] + 1), + ] def test_lower_bound(self) -> None: group_ids = self._seed(3, self.HASH) From 3d5e6fbec77df8269aea4da1fbee9dfe9f6e8eb4 Mon Sep 17 00:00:00 2001 From: Kyle Consalus Date: Tue, 22 Sep 2026 13:34:14 -0700 Subject: [PATCH 3/5] fix test --- tests/sentry/issues/derived/test_tasks.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/sentry/issues/derived/test_tasks.py b/tests/sentry/issues/derived/test_tasks.py index 030346a3c490..bfec2f437b66 100644 --- a/tests/sentry/issues/derived/test_tasks.py +++ b/tests/sentry/issues/derived/test_tasks.py @@ -1435,7 +1435,8 @@ def test_query_has_statement_timeout(self) -> None: with patch("sentry.issues.derived.tasks_util.statement_timeout") as timeout: group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) - assert timeout.call_args.args[1] == timedelta(seconds=50) + query_timeout = timeout.call_args.args[1] + assert timedelta(0) < query_timeout <= timedelta(seconds=40) def test_short_tail_is_one_range(self) -> None: null_ids = self._seed(3, None) From 71cbd6eae04cd3a59f58f80ffc32f65e74c0942b Mon Sep 17 00:00:00 2001 From: Kyle Consalus Date: Tue, 22 Sep 2026 13:40:46 -0700 Subject: [PATCH 4/5] simplify --- src/sentry/issues/derived/tasks.py | 4 +- src/sentry/issues/derived/tasks_util.py | 171 ++++++++++++++-------- tests/sentry/issues/derived/test_tasks.py | 62 +++++--- 3 files changed, 153 insertions(+), 84 deletions(-) diff --git a/src/sentry/issues/derived/tasks.py b/src/sentry/issues/derived/tasks.py index a5ac25e9ed41..b3b2788cfff0 100644 --- a/src/sentry/issues/derived/tasks.py +++ b/src/sentry/issues/derived/tasks.py @@ -591,8 +591,8 @@ def heal_stale_derived_data(**kwargs: object) -> None: try: range_result = group_id_ranges_for_hash( stale_hash, - chunk_size=batch_size, - max_chunks=remaining, + range_size=batch_size, + max_ranges=remaining, group_id_lower_bound=lower_bound, ) except OperationalError: diff --git a/src/sentry/issues/derived/tasks_util.py b/src/sentry/issues/derived/tasks_util.py index d0a96769a0c8..295c973c8099 100644 --- a/src/sentry/issues/derived/tasks_util.py +++ b/src/sentry/issues/derived/tasks_util.py @@ -3,12 +3,12 @@ import logging import random import time +from collections.abc import Sequence from dataclasses import dataclass from datetime import datetime, timedelta, timezone from math import ceil from typing import Protocol -from django.db import connections, router from django.db.models import Max, Min from django.db.utils import OperationalError @@ -22,9 +22,12 @@ _MAX_CHECK_GROUPS = 10_000 _GROUP_ID_RANGE_QUERY_TIMEOUT = timedelta(seconds=40) +# Use exact boundaries below this requested row count; estimate density above it. _MAX_EXACT_RANGE_ROWS = 10_000 +# Each probe reads this many matching IDs and extrapolates the following number of ranges. _RANGE_DENSITY_SAMPLE_SIZE = 100 _RANGES_PER_DENSITY_SAMPLE = 5 +# Bound total probe work for regions with large scheduling budgets. _MAX_RANGE_DENSITY_SAMPLES = 200 @@ -117,35 +120,84 @@ def _pick_random_fresh_group_ranges( return ranges +def _exact_group_id_ranges( + group_ids: Sequence[int], *, range_size: int, max_ranges: int +) -> list[tuple[int, int]]: + """Build exact ranges from ordered IDs, using an optional lookahead ID as the final end.""" + if not group_ids: + return [] + + requested_rows = range_size * max_ranges + starts = group_ids[:requested_rows:range_size] + last_end = group_ids[requested_rows] if len(group_ids) > requested_rows else group_ids[-1] + 1 + return list(zip(starts, [*starts[1:], last_end])) + + +def _estimate_group_id_ranges( + density_sample: Sequence[int], *, range_size: int, range_count: int +) -> list[tuple[int, int]]: + """Build contiguous ranges sized from the density of an ordered ID sample.""" + sample_width = density_sample[-1] - density_sample[0] + 1 + estimated_width = ceil(sample_width * range_size / len(density_sample)) + first_group_id = density_sample[0] + return [ + ( + first_group_id + index * estimated_width, + first_group_id + (index + 1) * estimated_width, + ) + for index in range(range_count) + ] + + def group_id_ranges_for_hash( - pipeline_hash: str | None, *, chunk_size: int, max_chunks: int, group_id_lower_bound: int = 0 + pipeline_hash: str | None, *, range_size: int, max_ranges: int, group_id_lower_bound: int = 0 ) -> GroupIdRangeResult: """Estimate ranges covering GroupDerivedData rows with a pipeline_hash. - Returns at most max_chunks of ascending disjoint [start, end) ranges, each - targeting chunk_size group IDs. Small requests use exact boundaries. Larger - requests sample local density once per five ranges, so estimate quality stays - stable as a region increases its scheduling budget. An over-dense range is - split by the worker, and an under-dense one costs only a scheduling slot. - ``drained`` is true only when a valid query found no rows at or above - ``group_id_lower_bound``. + Returns at most max_ranges of ascending disjoint [start, end) ranges, each + targeting range_size group IDs. ``drained`` is true only when a valid query + found no rows at or above ``group_id_lower_bound``. + + Imagine a sequence of GroupDerivedData rows with some stale (s): + + s.......s.......s.......s...........ssssssss.... + 0 8 16 24 36..43 + + We can precisely query for ranges with an equal number of stale rows, but that requires + us to do a full index scan of those rows, and if we've been mutating, that can get + surprisingly slow, especially since we'd like to be able to divvy out 100s of thousands + of rows. Our range processing task also filters and is tolerant of variation, so instead + of trying to be exact, we approximate. + + Simply cutting the ID span into N equal ranges can be rough. Asking for 3 above gives + [0,16) [16,32) [32,48), holding 2, 2, and 8 stale rows: one range has two thirds of the + work, and that's bad for our goal of great throughput. + + Instead, we set a budget of how much we'd like to do, and take incremental 'core samples', + using each to size the ranges that follow it. Sampling 4 rows at a time and targeting 4 + rows per range, the first sample reads 0, 8, 16, 24, four rows spanning 25 IDs, so we emit + [0,25). The next sample resumes there, skips the empty gap entirely, and lands on + 36, 37, 38, 39, four rows spanning 4 IDs, so we emit [36,40). The last sample sees the end + of the data and falls back to exact boundaries, [40,44). + + That's 4, 4, 4 instead of 2, 2, 8, for 3 small queries instead of a full scan. It's still + an estimate: an over-dense range is split by the worker, an under-dense one costs only a + scheduling slot. In production a sample is _RANGE_DENSITY_SAMPLE_SIZE rows and sizes + _RANGES_PER_DENSITY_SAMPLE ranges, rather than one. Raises ``OperationalError`` if all density probes together exceed the query budget. Callers must not interpret that as a drained hash. """ - if chunk_size <= 0 or max_chunks <= 0: + if range_size <= 0 or max_ranges <= 0: return GroupIdRangeResult(ranges=[], drained=False) - hash_predicate = "pipeline_hash IS NULL" if pipeline_hash is None else "pipeline_hash = %s" - sql = f""" - SELECT group_id - FROM {GroupDerivedData._meta.db_table} - WHERE {hash_predicate} AND group_id >= %s - ORDER BY pipeline_hash, group_id - LIMIT %s - """ - - using = router.db_for_read(GroupDerivedData) + matching_group_ids = ( + GroupDerivedData.objects.filter(pipeline_hash=pipeline_hash) + .order_by("pipeline_hash", "group_id") + .values_list("group_id", flat=True) + ) + using = matching_group_ids.db + matching_group_ids = matching_group_ids.using(using) query_deadline = time.monotonic() + _GROUP_ID_RANGE_QUERY_TIMEOUT.total_seconds() with metrics.timer("issues.derived.group_id_range_query"): @@ -155,65 +207,60 @@ def fetch_group_ids(start: int, limit: int) -> list[int]: if remaining_seconds <= 0.001: raise OperationalError("group ID range query budget exceeded") - params: list[str | int] = [] if pipeline_hash is None else [pipeline_hash] - params += [start, limit] - with ( - statement_timeout(using, timedelta(seconds=remaining_seconds)), - connections[using].cursor() as db_cursor, - ): - db_cursor.execute(sql, params) - return [row[0] for row in db_cursor.fetchall()] + with statement_timeout(using, timedelta(seconds=remaining_seconds)): + return list(matching_group_ids.filter(group_id__gte=start)[:limit]) - requested_rows = chunk_size * max_chunks + requested_rows = range_size * max_ranges if requested_rows <= _MAX_EXACT_RANGE_ROWS: group_ids = fetch_group_ids(group_id_lower_bound, requested_rows + 1) if not group_ids: return GroupIdRangeResult(ranges=[], drained=True) - starts = group_ids[:requested_rows:chunk_size] - if len(group_ids) > requested_rows: - ends = starts[1:] + [group_ids[requested_rows]] - else: - ends = starts[1:] + [group_ids[-1] + 1] - return GroupIdRangeResult(ranges=list(zip(starts, ends)), drained=False) + return GroupIdRangeResult( + ranges=_exact_group_id_ranges( + group_ids, + range_size=range_size, + max_ranges=max_ranges, + ), + drained=False, + ) - ranges: list[tuple[int, int]] = [] + result_ranges: list[tuple[int, int]] = [] next_group_id = group_id_lower_bound - sample_limit = _RANGE_DENSITY_SAMPLE_SIZE + 1 - sample_count = min( - ceil(max_chunks / _RANGES_PER_DENSITY_SAMPLE), + density_sample_count = min( + ceil(max_ranges / _RANGES_PER_DENSITY_SAMPLE), _MAX_RANGE_DENSITY_SAMPLES, ) - while len(ranges) < max_chunks: - sampled_group_ids = fetch_group_ids(next_group_id, sample_limit) + ranges_per_sample, samples_with_extra_range = divmod(max_ranges, density_sample_count) + # Ranges probably don't divide equally by samples, so we try to distribute the remainder + # cleanly. + range_counts = [ranges_per_sample + 1] * samples_with_extra_range + [ranges_per_sample] * ( + density_sample_count - samples_with_extra_range + ) + for ranges_for_sample in range_counts: + sampled_group_ids = fetch_group_ids(next_group_id, _RANGE_DENSITY_SAMPLE_SIZE + 1) if not sampled_group_ids: - return GroupIdRangeResult(ranges=ranges, drained=not ranges) - + return GroupIdRangeResult(ranges=result_ranges, drained=not result_ranges) if len(sampled_group_ids) <= _RANGE_DENSITY_SAMPLE_SIZE: - starts = sampled_group_ids[::chunk_size] + # We requested one extra, so this means we've got all the data and can exit. + starts = sampled_group_ids[::range_size] ends = starts[1:] + [sampled_group_ids[-1] + 1] - ranges.extend(zip(starts, ends)) + result_ranges.extend(zip(starts, ends)) break - density_sample = sampled_group_ids[:_RANGE_DENSITY_SAMPLE_SIZE] - sample_width = density_sample[-1] - density_sample[0] + 1 - estimated_chunk_width = max( - 1, - ceil(sample_width * chunk_size / len(density_sample)), + # trim the over-query + density_sample = sampled_group_ids[:-1] + # generate range_for_sample ranges based on density_sample. + estimated_ranges = _estimate_group_id_ranges( + density_sample, + range_size=range_size, + range_count=ranges_for_sample, ) - samples_left = min(sample_count, max_chunks - len(ranges)) - chunks_for_sample = ceil((max_chunks - len(ranges)) / samples_left) - - range_start = density_sample[0] - for _ in range(chunks_for_sample): - range_end = range_start + estimated_chunk_width - ranges.append((range_start, range_end)) - range_start = range_end - next_group_id = range_start - sample_count -= 1 - - return GroupIdRangeResult(ranges=ranges[:max_chunks], drained=False) + result_ranges.extend(estimated_ranges) + next_group_id = estimated_ranges[-1][1] + + return GroupIdRangeResult(ranges=result_ranges[:max_ranges], drained=False) def _resume_check_id( diff --git a/tests/sentry/issues/derived/test_tasks.py b/tests/sentry/issues/derived/test_tasks.py index bfec2f437b66..cd0cd82a4eca 100644 --- a/tests/sentry/issues/derived/test_tasks.py +++ b/tests/sentry/issues/derived/test_tasks.py @@ -31,6 +31,8 @@ from sentry.issues.derived.tasks_util import ( GroupIdRangeResult, SpawnState, + _estimate_group_id_ranges, + _exact_group_id_ranges, _pick_random_fresh_group_ranges, group_id_ranges_for_hash, ) @@ -829,8 +831,8 @@ def test_null_starts_at_zero_without_rewriting_state(self) -> None: mock_ranges.assert_called_once_with( None, - chunk_size=500, - max_chunks=1, + range_size=500, + max_ranges=1, group_id_lower_bound=0, ) mock_save.assert_not_called() @@ -1410,6 +1412,26 @@ def test_caps_total_groups(self) -> None: assert result == [(group_ids[0], group_ids[1] + 1)] +class TestGroupIdRangeMath: + def test_exact_ranges_use_lookahead_to_close_final_range(self) -> None: + assert _exact_group_id_ranges([10, 20, 30, 40, 50], range_size=2, max_ranges=2) == [ + (10, 30), + (30, 50), + ] + + def test_exact_ranges_close_short_tail_after_last_group(self) -> None: + assert _exact_group_id_ranges([10, 20, 30], range_size=2, max_ranges=5) == [ + (10, 30), + (30, 31), + ] + + def test_estimated_ranges_use_sample_density(self) -> None: + assert _estimate_group_id_ranges([10, 20, 31], range_size=2, range_count=2) == [ + (10, 25), + (25, 40), + ] + + @with_feature("projects:issue-action-log-write-to-db") class GroupIdRangesForHashTest(DerivedDataTaskTestBase): HASH = "a" * 16 @@ -1425,15 +1447,15 @@ def test_no_matching_rows(self) -> None: self._seed(2, self.OTHER_HASH) assert group_id_ranges_for_hash( - self.HASH, chunk_size=2, max_chunks=5 + self.HASH, range_size=2, max_ranges=5 ) == GroupIdRangeResult(ranges=[], drained=True) - assert group_id_ranges_for_hash(None, chunk_size=2, max_chunks=5) == GroupIdRangeResult( + assert group_id_ranges_for_hash(None, range_size=2, max_ranges=5) == GroupIdRangeResult( ranges=[], drained=True ) def test_query_has_statement_timeout(self) -> None: with patch("sentry.issues.derived.tasks_util.statement_timeout") as timeout: - group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) + group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5) query_timeout = timeout.call_args.args[1] assert timedelta(0) < query_timeout <= timedelta(seconds=40) @@ -1442,37 +1464,37 @@ def test_short_tail_is_one_range(self) -> None: null_ids = self._seed(3, None) hash_ids = self._seed(3, self.HASH) - # Fewer rows than chunk_size, for both the NULL and the concrete-hash + # Fewer rows than range_size, for both the NULL and the concrete-hash # predicate, and neither picks up the other's rows. - assert group_id_ranges_for_hash(None, chunk_size=10, max_chunks=5).ranges == [ + assert group_id_ranges_for_hash(None, range_size=10, max_ranges=5).ranges == [ (null_ids[0], null_ids[-1] + 1) ] - assert group_id_ranges_for_hash(self.HASH, chunk_size=10, max_chunks=5).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=10, max_ranges=5).ranges == [ (hash_ids[0], hash_ids[-1] + 1) ] def test_exact_chunk_boundaries(self) -> None: group_ids = self._seed(5, self.HASH) - assert group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5).ranges == [ (group_ids[0], group_ids[2]), (group_ids[2], group_ids[4]), (group_ids[4], group_ids[4] + 1), ] - def test_truncates_to_max_chunks(self) -> None: + def test_truncates_to_max_ranges(self) -> None: group_ids = self._seed(5, self.HASH) # The 5th group is left out rather than folded into an oversized last range. - assert group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=2).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=2).ranges == [ (group_ids[0], group_ids[2]), (group_ids[2], group_ids[4]), ] - def test_truncates_full_tail_to_max_chunks(self) -> None: + def test_truncates_full_tail_to_max_ranges(self) -> None: group_ids = self._seed(6, self.HASH) - assert group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=2).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=2).ranges == [ (group_ids[0], group_ids[2]), (group_ids[2], group_ids[4]), ] @@ -1481,10 +1503,10 @@ def test_invalid_chunking(self) -> None: self._seed(2, self.HASH) assert group_id_ranges_for_hash( - self.HASH, chunk_size=0, max_chunks=5 + self.HASH, range_size=0, max_ranges=5 ) == GroupIdRangeResult(ranges=[], drained=False) assert group_id_ranges_for_hash( - self.HASH, chunk_size=2, max_chunks=0 + self.HASH, range_size=2, max_ranges=0 ) == GroupIdRangeResult(ranges=[], drained=False) def test_density_probes_share_one_query_budget(self) -> None: @@ -1496,7 +1518,7 @@ def test_density_probes_share_one_query_budget(self) -> None: ), pytest.raises(OperationalError, match="query budget exceeded"), ): - group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) + group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5) def test_samples_local_density_for_approximate_ranges(self) -> None: groups = self.create_unprocessed_groups(81) @@ -1514,7 +1536,7 @@ def test_samples_local_density_for_approximate_ranges(self) -> None: patch("sentry.issues.derived.tasks_util._RANGES_PER_DENSITY_SAMPLE", 2), patch("sentry.issues.derived.tasks_util._MAX_RANGE_DENSITY_SAMPLES", 2), ): - result = group_id_ranges_for_hash(self.HASH, chunk_size=4, max_chunks=4) + result = group_id_ranges_for_hash(self.HASH, range_size=4, max_ranges=4) assert len(result.ranges) == 4 assert result.ranges == sorted(result.ranges) @@ -1530,7 +1552,7 @@ def test_low_volume_request_uses_exact_boundaries(self) -> None: patch("sentry.issues.derived.tasks_util._MAX_EXACT_RANGE_ROWS", 10), patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2), ): - result = group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) + result = group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5) assert result.ranges == [ (group_ids[0], group_ids[2]), @@ -1543,8 +1565,8 @@ def test_lower_bound(self) -> None: result = group_id_ranges_for_hash( self.HASH, - chunk_size=2, - max_chunks=5, + range_size=2, + max_ranges=5, group_id_lower_bound=group_ids[1], ) From 847ec9986c85ef037bb82bee59cae69bc006c4fc Mon Sep 17 00:00:00 2001 From: Kyle Consalus Date: Tue, 22 Sep 2026 16:25:35 -0700 Subject: [PATCH 5/5] fix comment --- src/sentry/issues/derived/tasks.py | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/sentry/issues/derived/tasks.py b/src/sentry/issues/derived/tasks.py index b3b2788cfff0..aa76c471cef0 100644 --- a/src/sentry/issues/derived/tasks.py +++ b/src/sentry/issues/derived/tasks.py @@ -967,9 +967,9 @@ def regenerate_stale_derived_data_batch( # How close the scheduler's density estimate landed to reality. The worker # records at most batch_size rows here and reports denser ranges separately as - # ``range_overflow`` reschedules. Counts well below batch_size mean ranges are - # too wide and slots are being wasted; frequent overflow means they are too - # narrow. Use both signals to tune the density sampling constants. + # ``range_overflow`` reschedules. Counts well below batch_size mean ranges span + # too few IDs and scheduling slots are being wasted; frequent overflow means they + # span too many. Use both signals to tune the density sampling constants. metrics.distribution( "issues.derived.heal_range_rows_found", len(group_ids),