diff --git a/src/sentry/conf/server.py b/src/sentry/conf/server.py index a3eaba4b7bbf..889fa8639a7e 100644 --- a/src/sentry/conf/server.py +++ b/src/sentry/conf/server.py @@ -2900,6 +2900,9 @@ def custom_parameter_sort(parameter: dict) -> tuple[str, int]: # How long reprocessing counters are kept in Redis before they expire. SENTRY_REPROCESSING_SYNC_TTL = 30 * 24 * 3600 # 30 days +# How long the reprocessing page claims are kept in Redis before they expire. +SENTRY_REPROCESSING_PAGE_CLAIM_TTL = 24 * 3600 # 1 day + # How many events to query for at once while paginating through an entire # issue. Note that this needs to be kept in sync with the time-limits on # `sentry.tasks.reprocessing2.reprocess_group`. That task is responsible for diff --git a/src/sentry/services/eventstore/reprocessing/base.py b/src/sentry/services/eventstore/reprocessing/base.py index b8f2d08d7d8f..02c90e07e825 100644 --- a/src/sentry/services/eventstore/reprocessing/base.py +++ b/src/sentry/services/eventstore/reprocessing/base.py @@ -1,6 +1,7 @@ from datetime import datetime from typing import Any, TypedDict +from sentry.utils.query import TaskBulkQueryState from sentry.utils.services import Service @@ -24,6 +25,7 @@ class ReprocessingStore(Service): "start_reprocessing", "get_pending", "get_progress", + "try_claim_page", ) def __init__(self, **options: Any) -> None: @@ -81,3 +83,13 @@ def get_pending(self, group_id: int) -> Any: def get_progress(self, group_id: int) -> ReprocessingInfo | None: raise NotImplementedError() + + def try_claim_page( + self, + project_id: int, + group_id: int, + new_group_id: int, + state: TaskBulkQueryState | None, + claimant: str, + ) -> bool: + raise NotImplementedError() diff --git a/src/sentry/services/eventstore/reprocessing/redis.py b/src/sentry/services/eventstore/reprocessing/redis.py index b70b97f6655a..95149f906c08 100644 --- a/src/sentry/services/eventstore/reprocessing/redis.py +++ b/src/sentry/services/eventstore/reprocessing/redis.py @@ -7,6 +7,7 @@ from django.conf import settings from sentry.utils.dates import to_datetime +from sentry.utils.query import TaskBulkQueryState from sentry.utils.redis import redis_clusters from .base import ReprocessingInfo, ReprocessingStore @@ -28,6 +29,12 @@ def _get_remaining_key(project_id: int, group_id: int) -> str: return f"re2:remaining:{{{project_id}:{group_id}}}" +def _get_page_claim_key( + project_id: int, group_id: int, new_group_id: int, timestamp: str, event_id: str +) -> str: + return f"re2:pageclaim:{project_id}:{group_id}:{new_group_id}:{timestamp}:{event_id}" + + class RedisReprocessingStore(ReprocessingStore): def __init__(self, **options: dict[str, Any]) -> None: cluster = options.pop("cluster", "default") @@ -180,3 +187,18 @@ def get_progress(self, group_id: int) -> ReprocessingInfo | None: if info is None: return None return orjson.loads(info) + + def try_claim_page( + self, + project_id: int, + group_id: int, + new_group_id: int, + state: TaskBulkQueryState | None, + claimant: str, + ) -> bool: + timestamp = state["timestamp"] if state is not None else "start" + event_id = state["event_id"] if state is not None else "start" + key = _get_page_claim_key(project_id, group_id, new_group_id, timestamp, event_id) + if self.redis.set(key, claimant, nx=True, ex=settings.SENTRY_REPROCESSING_PAGE_CLAIM_TTL): + return True + return self.redis.get(key) == claimant diff --git a/src/sentry/tasks/reprocessing2.py b/src/sentry/tasks/reprocessing2.py index e9104c1f80d8..6e69057b088f 100644 --- a/src/sentry/tasks/reprocessing2.py +++ b/src/sentry/tasks/reprocessing2.py @@ -20,6 +20,7 @@ from sentry.search.eap.occurrences.query_utils import build_group_id_in_filter from sentry.services import eventstore from sentry.services.eventstore.models import Event +from sentry.services.eventstore.reprocessing import reprocessing_store from sentry.silo.base import SiloMode from sentry.tasks.base import instrumented_task from sentry.tasks.process_buffer import buffer_incr @@ -109,6 +110,33 @@ def reprocess_group( assert new_group_id is not None + # INC-2445 follow-up: To the best of our knowledge we still have quite some `reprocess_group` tasks running in parallel, + # this logic is intended to cull all but one. This is temporary to recover from a bad state and should be dead code after that. + if activation_id: + try: + page_owned = reprocessing_store.try_claim_page( + project_id=project_id, + group_id=group_id, + new_group_id=new_group_id, + state=query_state, + claimant=activation_id, + ) + except Exception: + logger.warning("reprocessing2.page_claim.error", exc_info=True) + page_owned = True + if not page_owned: + logger.info( + "reprocessing2.page_claim.culled", + extra={ + "project_id": project_id, + "group_id": group_id, + "new_group_id": new_group_id, + "query_state": query_state, + "activation_id": activation_id, + }, + ) + return + query_state, events = task_run_batch_query( filter=eventstore.Filter(project_ids=[project_id], group_ids=[group_id]), batch_size=settings.SENTRY_REPROCESSING_PAGE_SIZE, diff --git a/tests/sentry/services/eventstore/processing/test_redis_cluster.py b/tests/sentry/services/eventstore/processing/test_redis_cluster.py index 8ab0626ec6fd..ece2df276fef 100644 --- a/tests/sentry/services/eventstore/processing/test_redis_cluster.py +++ b/tests/sentry/services/eventstore/processing/test_redis_cluster.py @@ -2,6 +2,7 @@ from sentry.services.eventstore.reprocessing.redis import RedisReprocessingStore from sentry.testutils.helpers.redis import use_redis_cluster +from sentry.utils.query import TaskBulkQueryState @use_redis_cluster() @@ -20,3 +21,22 @@ def test_mark_event_reprocessed() -> None: assert progress is not None assert progress.get("syncCount") == 10 assert progress.get("totalEvents") == 20 + + +@use_redis_cluster() +def test_try_claim_page() -> None: + store = RedisReprocessingStore() + project_id = 1 + group_id = 2 + new_group_id = 3 + state: TaskBulkQueryState = {"timestamp": "2026-08-04T06:10:59+00:00", "event_id": "42"} + + # First claim is ok, and reclaiming from same claimant is a NOOP. + assert store.try_claim_page(project_id, group_id, new_group_id, state, claimant="A") + assert store.try_claim_page(project_id, group_id, new_group_id, state, claimant="A") + + # Claiming from another claimant should not work. + assert not store.try_claim_page(project_id, group_id, new_group_id, state, claimant="B") + + # Different reprocessing run but same state should be unaffected. + assert store.try_claim_page(project_id, 4, 5, state, claimant="B")