Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 42 additions & 4 deletions src/sentry/issues/derived/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -581,14 +581,18 @@ 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:
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:
Expand All @@ -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(
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -944,14 +951,30 @@ 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
group_id__gte=group_id_start,
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
Comment thread
kcons marked this conversation as resolved.
if range_overflow:
group_ids = group_ids[:batch_size]

# 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 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),
sample_rate=1.0,
tags={"hash_kind": "null" if target_hash is None else "stale"},
)

result = build_and_promote_batch(
Expand Down Expand Up @@ -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,
Expand Down
223 changes: 146 additions & 77 deletions src/sentry/issues/derived/tasks_util.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,12 +2,15 @@

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

from sentry.issues.derived.check import CheckFailure, CheckId, CheckInvalidated, CheckResult
from sentry.issues.models.groupderiveddata import GroupDerivedData
Expand All @@ -18,11 +21,14 @@
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)
# 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


@dataclass(frozen=True)
Expand Down Expand Up @@ -114,84 +120,147 @@ 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:
"""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_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):

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``.
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)

# 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,
},
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"):

def fetch_group_ids(start: int, limit: int) -> list[int]:
remaining_seconds = query_deadline - time.monotonic()

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

nb: I have a helper for this sitting around, but didn't want to make it part of this pr. See #125142

if remaining_seconds <= 0.001:
raise OperationalError("group ID range query budget exceeded")

with statement_timeout(using, timedelta(seconds=remaining_seconds)):
return list(matching_group_ids.filter(group_id__gte=start)[:limit])

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)

return GroupIdRangeResult(
ranges=_exact_group_id_ranges(
group_ids,
range_size=range_size,
max_ranges=max_ranges,
),
drained=False,
)

result_ranges: list[tuple[int, int]] = []
next_group_id = group_id_lower_bound
density_sample_count = min(
ceil(max_ranges / _RANGES_PER_DENSITY_SAMPLE),
_MAX_RANGE_DENSITY_SAMPLES,
)
# 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
"""
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)
with (
statement_timeout(using, _GROUP_ID_RANGE_TIMEOUT),
metrics.timer("issues.derived.group_id_range_query"),
connections[using].cursor() as 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)
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=result_ranges, drained=not result_ranges)
if len(sampled_group_ids) <= _RANGE_DENSITY_SAMPLE_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]
result_ranges.extend(zip(starts, ends))
break
Comment thread
kcons marked this conversation as resolved.

# 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,
)
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(
Expand Down
Loading
Loading