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
6 changes: 6 additions & 0 deletions src/sentry/options/defaults.py
Original file line number Diff line number Diff line change
Expand Up @@ -2991,6 +2991,12 @@
default=False,
flags=FLAG_PRIORITIZE_DISK | FLAG_AUTOMATOR_MODIFIABLE,
)
register(
"spans.buffer.process-segments-task-rollout-rate",
type=Float,
default=0.0,
flags=FLAG_PRIORITIZE_DISK | FLAG_AUTOMATOR_MODIFIABLE,
)

# List of trace_ids to enable debug logging for. Empty = debug off.
# When set, logs detailed metrics about zunionstore set sizes, key existence, and trace structure.
Expand Down
33 changes: 27 additions & 6 deletions src/sentry/spans/consumers/process/flusher.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from collections.abc import Callable, Mapping
from concurrent.futures import Future
from functools import partial
from typing import Any

import orjson
import sentry_sdk
Expand All @@ -21,8 +22,10 @@
from sentry.conf.types.kafka_definition import Topic
from sentry.constants import DataCategory
from sentry.models.project import Project
from sentry.options.rollout import in_random_rollout
from sentry.processing.backpressure.memory import ServiceMemory
from sentry.spans.buffer import SpansBuffer
from sentry.spans.consumers.process_segments.tasks import process_segment_task
from sentry.utils import metrics
from sentry.utils.arroyo import run_with_initialized_sentry
from sentry.utils.kafka_config import get_kafka_producer_cluster_options, get_topic_definition
Expand Down Expand Up @@ -280,7 +283,7 @@ def main(
logger.info("Flusher process started for shards %s", shard_tag)

try:
producer_futures = []
producer_futures: list[tuple[int, Any, int]] = []

if produce_to_pipe is not None:

Expand Down Expand Up @@ -351,18 +354,36 @@ def produce(project_id: int, payload: KafkaPayload, dropped: int) -> None:
},
)

produce_to_process_segment_task = (
produce_to_pipe is None
and in_random_rollout("spans.buffer.process-segments-task-rollout-rate")
)

for message in flushed_segment.to_messages():
kafka_payload = KafkaPayload(None, orjson.dumps(message), [])
metrics.timing(
"spans.buffer.segment_size_bytes",
len(kafka_payload.value),
tags={"shard": shard_tag},
)
produce(
flushed_segment.project_id,
kafka_payload,
len(message["spans"]),
)
if produce_to_process_segment_task:
task_produce_future = process_segment_task.apply_async_with_future(
args=[kafka_payload.value]
)
if task_produce_future is not None:
producer_futures.append(
(
flushed_segment.project_id,
task_produce_future,
len(message["spans"]),
)
)
else:
produce(
flushed_segment.project_id,
kafka_payload,
len(message["spans"]),
)

with metrics.timer("spans.buffer.flusher.wait_produce", tags={"shards": shard_tag}):
for project_id, future, dropped in producer_futures:
Expand Down
2 changes: 1 addition & 1 deletion src/sentry/taskworker/namespaces.py
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@

spans_process_segments_tasks = app.taskregistry.create_namespace(
"spans.process_segments",
app_feature="transactions",
app_feature="spans",
)

ingest_attachments_tasks = app.taskregistry.create_namespace(
Expand Down
2 changes: 1 addition & 1 deletion src/sentry/taskworker/silolimiter.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ def handle(*args: P.args, **kwargs: P.kwargs) -> Any:

def __call__(self, decorated_task: Task[P, R]) -> Task[P, R]:
# Replace the sentry.taskworker.Task interface used to schedule tasks.
replacements = {"delay", "apply_async"}
replacements = {"delay", "apply_async", "apply_async_with_future"}
for attr_name in replacements:
task_attr = getattr(decorated_task, attr_name)
if callable(task_attr):
Expand Down
147 changes: 132 additions & 15 deletions tests/sentry/spans/consumers/process/test_flusher.py
Original file line number Diff line number Diff line change
@@ -1,11 +1,13 @@
import time
from concurrent.futures import Future
from time import sleep
from types import SimpleNamespace
from typing import Any
from unittest import mock

import orjson
import pytest
from arroyo.backends.kafka import KafkaPayload
from arroyo.processing.strategies.noop import Noop
from django.test import override_settings

Expand All @@ -21,6 +23,32 @@ def _payload(span_id: str) -> bytes:
return orjson.dumps({"span_id": span_id})


def _buffer_with_segment(
*,
project_id: int = 999_002,
trace_id: str = "9" * 32,
span_id: str = "b" * 16,
slice_id: int = 999_002,
) -> SpansBuffer:
buffer = SpansBuffer(assigned_shards=[0], slice_id=slice_id)
buffer.process_spans(
[
Span(
payload=_payload(span_id),
trace_id=trace_id,
span_id=span_id,
parent_span_id=None,
segment_id=None,
is_segment_span=True,
project_id=project_id,
partition=0,
)
],
now=0,
)
return buffer


def _blocking_main_for_join_test(
buffer, shards, stopped, current_drift, backpressure_since, healthy_since, produce_to_pipe
):
Expand All @@ -39,21 +67,8 @@ def test_flusher_logs_flushed_segments() -> None:
slice_id = 999_002
segment_key = f"span-buf:s:{{{project_id}:{trace_id}}}:{span_id}".encode()
queue_key = f"span-buf:q:{slice_id}-0".encode()
buffer = SpansBuffer(assigned_shards=[0], slice_id=slice_id)
buffer.process_spans(
[
Span(
payload=_payload(span_id),
trace_id=trace_id,
span_id=span_id,
parent_span_id=None,
segment_id=None,
is_segment_span=True,
project_id=project_id,
partition=0,
)
],
now=0,
buffer = _buffer_with_segment(
project_id=project_id, trace_id=trace_id, span_id=span_id, slice_id=slice_id
)
stopped = SimpleNamespace(value=0)
current_drift = SimpleNamespace(value=0)
Expand Down Expand Up @@ -93,6 +108,108 @@ def produce_to_pipe(project_id: int, payload: Any, dropped: int) -> None:
assert not buffer.client.keys(f"*{project_id}:{trace_id}*")


@override_options({**DEFAULT_OPTIONS, "spans.buffer.process-segments-task-rollout-rate": 0.0})
def test_flusher_produces_flushed_segments_to_buffered_segments_topic() -> None:
span_id = "b" * 16
buffer = _buffer_with_segment(span_id=span_id)
stopped = SimpleNamespace(value=0)
current_drift = SimpleNamespace(value=0)
backpressure_since = SimpleNamespace(value=0)
healthy_since = SimpleNamespace(value=0)
produced_payloads: list[KafkaPayload] = []
producer_future: Future[None] = Future()
producer_future.set_result(None)
producer_manager = mock.Mock()

def produce(payload: KafkaPayload) -> Future[None]:
produced_payloads.append(payload)
stopped.value = 1
return producer_future

producer_manager.produce.side_effect = produce

with (
mock.patch(
"sentry.spans.consumers.process.flusher.MultiProducer",
return_value=producer_manager,
) as mock_multi_producer,
mock.patch(
"sentry.spans.consumers.process.flusher.process_segment_task.apply_async_with_future"
) as mock_apply_async_with_future,
):
SpanFlusher.main(
buffer,
shards=[0],
stopped=stopped,
current_drift=current_drift,
backpressure_since=backpressure_since,
healthy_since=healthy_since,
produce_to_pipe=None,
)

mock_multi_producer.assert_called_once_with(Topic.BUFFERED_SEGMENTS)
mock_apply_async_with_future.assert_not_called()
assert len(produced_payloads) == 1
assert orjson.loads(produced_payloads[0].value)["spans"][0]["span_id"] == span_id


@override_options({**DEFAULT_OPTIONS, "spans.buffer.process-segments-task-rollout-rate": 1.0})
def test_flusher_produces_flushed_segments_to_process_segment_task() -> None:
project_id = 999_002
trace_id = "9" * 32
span_id = "c" * 16
slice_id = 999_002
segment_key = f"span-buf:s:{{{project_id}:{trace_id}}}:{span_id}".encode()
queue_key = f"span-buf:q:{slice_id}-0".encode()
buffer = _buffer_with_segment(
project_id=project_id, trace_id=trace_id, span_id=span_id, slice_id=slice_id
)
stopped = SimpleNamespace(value=0)
current_drift = SimpleNamespace(value=0)
backpressure_since = SimpleNamespace(value=0)
healthy_since = SimpleNamespace(value=0)
producer_manager = mock.Mock()
task_future = mock.Mock()

def wait_for_task_delivery() -> None:
assert buffer.client.zscore(queue_key, segment_key) is not None

task_future.result.side_effect = wait_for_task_delivery

def apply_async_with_future(*args: Any, **kwargs: Any) -> Any:
stopped.value = 1
return task_future

with (
mock.patch(
"sentry.spans.consumers.process.flusher.MultiProducer",
return_value=producer_manager,
) as mock_multi_producer,
mock.patch(
"sentry.spans.consumers.process.flusher.process_segment_task.apply_async_with_future",
side_effect=apply_async_with_future,
) as mock_apply_async_with_future,
):
SpanFlusher.main(
buffer,
shards=[0],
stopped=stopped,
current_drift=current_drift,
backpressure_since=backpressure_since,
healthy_since=healthy_since,
produce_to_pipe=None,
)

mock_multi_producer.assert_called_once_with(Topic.BUFFERED_SEGMENTS)
producer_manager.produce.assert_not_called()
mock_apply_async_with_future.assert_called_once()
task_args = mock_apply_async_with_future.call_args.kwargs["args"]
assert len(task_args) == 1
assert orjson.loads(task_args[0])["spans"][0]["span_id"] == span_id
task_future.result.assert_called_once_with()
assert buffer.client.zscore(queue_key, segment_key) is None


@override_options({**DEFAULT_OPTIONS, "spans.buffer.max-flush-segments": 1})
def test_backpressure() -> None:
# Flush very aggressively to make join() faster
Expand Down
2 changes: 1 addition & 1 deletion tests/sentry/tasks/test_base.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,7 +93,7 @@ def test_task_silo_limit_call_monolith() -> None:

@pytest.mark.parametrize(
"method_name",
("apply_async", "delay"),
("apply_async", "apply_async_with_future", "delay"),
)
@override_settings(SILO_MODE=SiloMode.CONTROL)
def test_task_silo_limit_task_methods(method_name: str) -> None:
Expand Down
Loading