From 386f2a83eba64d5989c86feca35a269684f001ee Mon Sep 17 00:00:00 2001 From: Yucheng He Date: Fri, 4 Sep 2026 17:05:14 -0700 Subject: [PATCH] fix(otel v2): build one destination processor per destination, not per racing span Building outside the cache lock meant a cold cache met by a burst of concurrent requests constructed an exporter per thread, kept one, and handed the rest to the drain, so a batch worker and a connection pool per losing thread sat in a queue two workers service. Build under the lock that reads the cache. Opening an exporter connects to nothing, so the lock is held for a constructor, once per destination, and the race it was avoiding stops existing. --- .../integrations/otel/plumbing/providers.py | 42 ++++++++--------- .../otel/test_otel_v2_destinations.py | 46 ++++++++++++++++--- 2 files changed, 59 insertions(+), 29 deletions(-) diff --git a/litellm/integrations/otel/plumbing/providers.py b/litellm/integrations/otel/plumbing/providers.py index c30e0e7e5b2..f70f07a8ef8 100644 --- a/litellm/integrations/otel/plumbing/providers.py +++ b/litellm/integrations/otel/plumbing/providers.py @@ -426,36 +426,34 @@ class TenantFanOutSpanProcessor(SpanProcessor): return False def _acquire(self, destination: "OtelDestination") -> SpanProcessor | None: - """The processor for ``destination``, marked busy until ``_release``.""" + """The processor for ``destination``, marked busy until ``_release``. + + The build happens under the same lock that reads the cache, so a cold cache + met by a burst of concurrent requests yields one exporter rather than one per + thread with all but the winner shed. Building an exporter opens no connection, + so the cost of holding the lock is a constructor, once per destination. + """ key: Final = destination.cache_key() with self._lock: if self._closed: return None - cached = self._processors.get(key) # rebind-ok: reassigned after the build below - if cached is not None: + if (cached := self._processors.get(key)) is not None: self._processors.move_to_end(key) - self._exporting[id(cached)] = self._exporting.get(id(cached), 0) + 1 - return cached + processor: Final = cached if cached is not None else self._build_locked(destination, key) + if processor is None: + return None + self._exporting[id(processor)] = self._exporting.get(id(processor), 0) + 1 + drained: Final = self._drainable_locked() + for shed in drained: + self._drain.submit(shed) + return processor + + def _build_locked(self, destination: "OtelDestination", key: object) -> SpanProcessor | None: built: Final = self._build(destination) if built is None: return None - with self._lock: - if self._closed: - # Shutdown ran while this one was being built, so it belongs to nobody. - _shutdown_quietly(built) - return None - existing: Final = self._processors.get(key) - if existing is not None: - # Another thread won the race; drop ours rather than leak its thread. - self._drain.submit(built) - self._exporting[id(existing)] = self._exporting.get(id(existing), 0) + 1 - return existing - self._processors[key] = built - self._exporting[id(built)] = 1 - self._retire_overflow_locked() - drained: Final = self._drainable_locked() - for processor in drained: - self._drain.submit(processor) + self._processors[key] = built + self._retire_overflow_locked() return built def _release(self, processor: SpanProcessor) -> None: diff --git a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py index f51dc27a690..d27ccedf18c 100644 --- a/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py +++ b/tests/test_litellm/integrations/otel/test_otel_v2_destinations.py @@ -1217,10 +1217,10 @@ class TestEvictionSafety: assert held.shutdown_calls == 1 - def test_a_processor_built_during_shutdown_is_not_left_in_a_cleared_cache(self): - """The build runs outside the lock, so shutdown can finish inside it. Inserting - afterwards leaves a live exporter, with its batch thread and its connection - pool, in a map nothing will read again.""" + def test_a_processor_built_while_shutdown_waits_is_still_closed(self): + """Shutdown cannot slip between the build and the insert, which would leave a + live exporter, with its batch thread and its connection pool, in a map nothing + will read again.""" import threading built = [] @@ -1230,7 +1230,7 @@ class TestEvictionSafety: built.append(self.Recording()) return built[-1] - fan_out = TenantFanOutSpanProcessor(processor_factory=slow) + fan_out = TenantFanOutSpanProcessor(processor_factory=slow, shutdown_drain_seconds=0.05) acquired = [] caller = threading.Thread(target=lambda: acquired.append(fan_out._acquire(self._dest(0)))) caller.start() @@ -1238,9 +1238,41 @@ class TestEvictionSafety: fan_out.shutdown() caller.join(timeout=10) - assert acquired == [None], "an exporter built after shutdown was handed out" + assert acquired == built, "the build shutdown waited out was thrown away" + + fan_out._release(built[0]) + self._settle(fan_out, built[0]) + + assert built[0].shutdown_calls == 1, "the exporter outlived the fan-out" assert fan_out._processors == {}, "an exporter was left in a cleared cache" - assert built[0].shutdown_calls == 1, "the exporter that lost the race was never closed" + + def test_a_cold_cache_met_by_a_burst_builds_one_processor_per_destination(self): + """Building outside the cache lock let every thread of the burst construct its + own exporter, each with a batch thread and a connection pool, and shed all but + one into the drain.""" + import threading + + built = [] + + def factory(_destination): + time.sleep(0.01) + built.append(self.Recording()) + return built[-1] + + fan_out = TenantFanOutSpanProcessor(processor_factory=factory) + ready = threading.Barrier(8) + + def acquire(): + ready.wait() + fan_out._release(fan_out._acquire(self._dest(0))) + + callers = [threading.Thread(target=acquire) for _ in range(8)] + for caller in callers: + caller.start() + for caller in callers: + caller.join(timeout=10) + + assert len(built) == 1, f"one destination, {len(built)} exporters built" def test_a_submit_racing_close_is_never_stranded_behind_the_sentinels(self): """A submit that read the closed state and then let ``close`` run queues its