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.
This commit is contained in:
Yucheng He 2026-09-04 17:05:14 -07:00
parent c366b91729
commit 386f2a83eb
2 changed files with 59 additions and 29 deletions

View file

@ -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:

View file

@ -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