mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-10 03:28:53 +00:00
fix(otel): bound shed destination processors waiting on a dead collector
This commit is contained in:
parent
4372884404
commit
093a9912c8
2 changed files with 75 additions and 2 deletions
|
|
@ -216,6 +216,12 @@ _MAX_CACHED_DESTINATION_PROCESSORS: Final = 32
|
|||
#: create by cycling its destination config.
|
||||
_DRAIN_WORKERS: Final = 2
|
||||
|
||||
#: Shed processors waiting to be closed before the fan-out stops building new ones.
|
||||
#: Each still owns a batch thread until its close returns, and a collector that never
|
||||
#: answers makes every close take the exporter's full timeout, so past this many the
|
||||
#: operator's exporter keeps the span instead (see ``deliverable``).
|
||||
_MAX_PENDING_DRAINS: Final = 64
|
||||
|
||||
#: How long ``shutdown`` waits for spans already being forwarded, so teardown closes
|
||||
#: no processor under one. Bounded: an exporter that never returns must not hold the
|
||||
#: proxy open.
|
||||
|
|
@ -247,10 +253,13 @@ class _DrainPool:
|
|||
self,
|
||||
workers: int = _DRAIN_WORKERS,
|
||||
pending: "queue.Queue[SpanProcessor | None] | None" = None,
|
||||
capacity: int = _MAX_PENDING_DRAINS,
|
||||
) -> None:
|
||||
self._workers: Final = workers
|
||||
self._capacity: Final = capacity
|
||||
self._lock: Final = threading.Lock()
|
||||
self._closed = False
|
||||
self._backlog = 0 # guarded by ``_lock``: submitted processors whose close has not returned
|
||||
self._pending: Final[queue.Queue[SpanProcessor | None]] = pending if pending is not None else queue.Queue()
|
||||
self._threads: Final = tuple(
|
||||
threading.Thread(target=self._drain_until_closed, daemon=True, name="litellm-otel-destination-drain")
|
||||
|
|
@ -274,6 +283,7 @@ class _DrainPool:
|
|||
"""
|
||||
with self._lock:
|
||||
if not self._closed:
|
||||
self._backlog += 1
|
||||
self._pending.put(processor)
|
||||
return
|
||||
threading.Thread(
|
||||
|
|
@ -283,6 +293,18 @@ class _DrainPool:
|
|||
name="litellm-otel-destination-drain-straggler",
|
||||
).start()
|
||||
|
||||
def saturated(self) -> bool:
|
||||
"""Whether enough closes are outstanding that building another processor must wait.
|
||||
|
||||
The workers close in order and each close blocks for as long as its exporter
|
||||
does, so a collector that stopped answering would otherwise turn every new
|
||||
destination into one more batch thread parked behind them, for as long as the
|
||||
tenants keep rotating. Holding the count here rather than reading the queue
|
||||
keeps the two processors a worker is mid-close on in the total.
|
||||
"""
|
||||
with self._lock:
|
||||
return self._backlog >= self._capacity
|
||||
|
||||
def close(self, timeout: float | None = None) -> None:
|
||||
"""Retire the workers once they have closed everything already queued.
|
||||
|
||||
|
|
@ -311,6 +333,8 @@ class _DrainPool:
|
|||
if processor is None:
|
||||
return
|
||||
_shutdown_quietly(processor)
|
||||
with self._lock:
|
||||
self._backlog -= 1
|
||||
|
||||
|
||||
_NO_ATTRIBUTES: Final[Mapping[str, AttributeValue]] = MappingProxyType({})
|
||||
|
|
@ -405,6 +429,7 @@ class TenantFanOutSpanProcessor(SpanProcessor):
|
|||
processor_factory: 'Callable[["OtelDestination"], SpanProcessor | None] | None' = None,
|
||||
shutdown_drain_seconds: float = _SHUTDOWN_DRAIN_SECONDS,
|
||||
operator_sinks: frozenset[_SinkKey] = frozenset(),
|
||||
pending_drains: int = _MAX_PENDING_DRAINS,
|
||||
) -> None:
|
||||
self._operator_sinks: Final = operator_sinks
|
||||
self._drain_seconds: Final = shutdown_drain_seconds
|
||||
|
|
@ -414,7 +439,7 @@ class TenantFanOutSpanProcessor(SpanProcessor):
|
|||
self._processors: OrderedDict[object, SpanProcessor] = OrderedDict() # mutable-ok: bounded LRU
|
||||
self._retired: OrderedDict[int, SpanProcessor] = OrderedDict() # mutable-ok: drains as exports finish
|
||||
self._exporting: dict[int, int] = {} # mutable-ok: per-processor in-flight export count
|
||||
self._drain: Final = _DrainPool()
|
||||
self._drain: Final = _DrainPool(capacity=pending_drains)
|
||||
|
||||
def on_start(self, span: SDKSpan, parent_context: Context | None = None) -> None:
|
||||
return None
|
||||
|
|
@ -544,6 +569,16 @@ class TenantFanOutSpanProcessor(SpanProcessor):
|
|||
return self._build_locked(destination, key)
|
||||
|
||||
def _build_locked(self, destination: "OtelDestination", key: object) -> SpanProcessor | None:
|
||||
"""Build and cache a processor for ``destination``, unless the drain is saturated.
|
||||
|
||||
Every build past the cache cap sheds one processor into the drain, so while the
|
||||
shed ones are stuck closing against a collector that stopped answering, a new
|
||||
destination is refused rather than parked behind them: ``deliverable`` then
|
||||
leaves its spans with the operator's exporter until the drain catches up.
|
||||
"""
|
||||
if self._drain.saturated():
|
||||
verbose_logger.debug("OTel V2 fan-out: drain saturated, not building for %s", destination.endpoint)
|
||||
return None
|
||||
built: Final = self._build(destination)
|
||||
if built is None:
|
||||
return None
|
||||
|
|
|
|||
|
|
@ -1676,7 +1676,45 @@ class TestEvictionSafety:
|
|||
release.set()
|
||||
self._settle(fan_out, built[0])
|
||||
|
||||
def test_the_drain_workers_do_not_hold_the_process_open(self):
|
||||
def test_a_saturated_drain_leaves_new_destinations_with_the_operator(self):
|
||||
"""A shed processor keeps its batch thread until its close returns, and against
|
||||
a collector that never answers every close waits out the exporter's timeout.
|
||||
Tenants rotating past the cache cap would otherwise queue one more processor,
|
||||
and one more thread, per request for as long as the outage lasts."""
|
||||
import threading
|
||||
|
||||
from litellm.integrations.otel.plumbing.providers import _MAX_CACHED_DESTINATION_PROCESSORS
|
||||
|
||||
release = threading.Event()
|
||||
|
||||
class Blocking(self.Recording):
|
||||
def shutdown(self):
|
||||
release.wait(timeout=10)
|
||||
super().shutdown()
|
||||
|
||||
built = []
|
||||
|
||||
def factory(_destination):
|
||||
built.append(Blocking())
|
||||
return built[-1]
|
||||
|
||||
fan_out = TenantFanOutSpanProcessor(processor_factory=factory, pending_drains=3)
|
||||
try:
|
||||
for index in range(_MAX_CACHED_DESTINATION_PROCESSORS + 40):
|
||||
processor = fan_out._acquire(self._dest(index))
|
||||
if processor is not None:
|
||||
fan_out._release(processor)
|
||||
|
||||
assert len(built) == _MAX_CACHED_DESTINATION_PROCESSORS + 3, "a processor per request during the outage"
|
||||
assert fan_out.deliverable((self._dest(999),)) == (), "the span would vanish instead of staying with the operator"
|
||||
finally:
|
||||
release.set()
|
||||
for _ in range(500):
|
||||
if not fan_out._drain.saturated():
|
||||
break
|
||||
time.sleep(0.02)
|
||||
|
||||
assert fan_out.deliverable((self._dest(999),)) == (self._dest(999),), "the fan-out never recovered"
|
||||
"""Python joins a ThreadPoolExecutor's workers at interpreter exit, so one
|
||||
unreachable tenant collector would hold the proxy open for its export
|
||||
timeout on the way down."""
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue