diff --git a/sentry_sdk/_span_batcher.py b/sentry_sdk/_span_batcher.py index c1c40044ee..2033cf7845 100644 --- a/sentry_sdk/_span_batcher.py +++ b/sentry_sdk/_span_batcher.py @@ -30,6 +30,7 @@ class SpanBatcher(Batcher["SpanJSON"]): GLOBAL_MAX_BEFORE_DROP = 10_000 MAX_BYTES_BEFORE_FLUSH = 5 * 1024 * 1024 # 5 MB + GLOBAL_MAX_BYTES_BEFORE_FLUSH = 25 * 1024 * 1024 # 25 MB FLUSH_WAIT_TIME = 5.0 @@ -50,6 +51,8 @@ def __init__( self._span_number: int = 0 self._running_size: dict[str, int] = defaultdict(lambda: 0) + self._total_running_size: int = 0 + self._capture_func = capture_func self._record_lost_func = record_lost_func self._running = True @@ -79,6 +82,8 @@ def _reset_thread_state(self) -> None: self._span_number = 0 self._running_size = defaultdict(lambda: 0) + self._total_running_size = 0 + self._running = True self._lock = threading.Lock() @@ -100,7 +105,7 @@ def _flush_loop(self) -> None: self._flush(only_pending=True) - if ( + if self._total_running_size >= self.GLOBAL_MAX_BYTES_BEFORE_FLUSH or ( time.monotonic() - self._last_full_flush >= self.FLUSH_WAIT_TIME + jitter ): @@ -137,7 +142,9 @@ def add(self, span: "SpanJSON") -> None: self._span_buffer[span["trace_id"]].append(span) self._span_number += 1 - self._running_size[span["trace_id"]] += self._estimate_size(span) + estimated_size = self._estimate_size(span) + self._running_size[span["trace_id"]] += estimated_size + self._total_running_size += estimated_size if ( len(self._span_buffer[span["trace_id"]]) >= self.MAX_BEFORE_FLUSH @@ -147,7 +154,9 @@ def add(self, span: "SpanJSON") -> None: self._pending_flush.add(span["trace_id"]) notify = True else: - notify = False + notify = ( + self._total_running_size >= self.GLOBAL_MAX_BYTES_BEFORE_FLUSH + ) if notify: self._flush_event.set() @@ -241,6 +250,7 @@ def _flush(self, only_pending: bool = False) -> None: self._span_number -= len(self._span_buffer[bucket_id]) del self._span_buffer[bucket_id] + self._total_running_size -= self._running_size[bucket_id] del self._running_size[bucket_id] for envelope in envelopes: diff --git a/tests/tracing/test_span_batcher.py b/tests/tracing/test_span_batcher.py index 507fef184b..fa34ab26e1 100644 --- a/tests/tracing/test_span_batcher.py +++ b/tests/tracing/test_span_batcher.py @@ -338,6 +338,69 @@ def test_weight_based_flushing_by_attribute_size( assert envelopes[0].items[0].payload.json["items"][1]["name"] == "big span" +def test_global_length_based_flushing(sentry_init, capture_items, monkeypatch): + """When the batcher reaches GLOBAL_MAX_BYTES_BEFORE_FLUSH, all buckets will be flushed.""" + # Limit of 2_000 is just above the size of a bare span. + monkeypatch.setattr(SpanBatcher, "GLOBAL_MAX_BYTES_BEFORE_FLUSH", 2_000) + # set the time-based flush limit to something huge so that it doesn't + # interfere + monkeypatch.setattr(SpanBatcher, "FLUSH_WAIT_TIME", 100000) + + sentry_init( + traces_sample_rate=1.0, + trace_lifecycle="stream", + ) + + items = capture_items("span") + + with sentry_sdk.traces.start_span(name="span"): + pass + + sentry_sdk.traces.new_trace() + with sentry_sdk.traces.start_span(name="span"): + pass + + time.sleep(0.1) + + assert len(items) == 2 + assert items[0].payload["name"] == "span" + + +def test_total_size_reset_after_length_based_flushing( + sentry_init, capture_items, monkeypatch +): + """Span is not flushed after a flush reduces the combined span size in bytes below the global limit.""" + # Limit of 2_000 is just above the size of a bare span. + monkeypatch.setattr(SpanBatcher, "GLOBAL_MAX_BYTES_BEFORE_FLUSH", 2_000) + # set the time-based flush limit to something huge so that it doesn't + # interfere + monkeypatch.setattr(SpanBatcher, "FLUSH_WAIT_TIME", 100000) + + sentry_init( + traces_sample_rate=1.0, + trace_lifecycle="stream", + ) + + items = capture_items("span") + + with sentry_sdk.traces.start_span(name="span"): + pass + + sentry_sdk.traces.new_trace() + with sentry_sdk.traces.start_span(name="span"): + pass + + time.sleep(0.1) + + with sentry_sdk.traces.start_span(name="span"): + pass + + time.sleep(0.1) + + assert len(items) == 2 + assert items[0].payload["name"] == "span" + + def test_bucket_recreated_after_flush(sentry_init, capture_envelopes, monkeypatch): """Spans for a trace that arrive after that trace's bucket was flushed land in a fresh bucket.""" monkeypatch.setattr(SpanBatcher, "MAX_BEFORE_FLUSH", 2) @@ -545,6 +608,8 @@ def test_span_batcher_lock_reset_in_child_after_fork(sentry_init): batcher._span_number = 1 batcher._running_size["test-trace-id"] = 42 + batcher._total_running_size = 42 + batcher._active.flag = True batcher._flush_event.set() batcher._running = False @@ -559,6 +624,7 @@ def test_span_batcher_lock_reset_in_child_after_fork(sentry_init): span_number_reset = batcher._span_number == 0 running_size_reset = len(batcher._running_size) == 0 + total_running_size_reset = batcher._total_running_size == 0 active_reset = not getattr(batcher._active, "flag", False) event_reset = not batcher._flush_event.is_set() @@ -572,6 +638,7 @@ def test_span_batcher_lock_reset_in_child_after_fork(sentry_init): and span_buffer_reset and span_number_reset and running_size_reset + and total_running_size_reset and active_reset and event_reset and running_reset