mirror of
https://github.com/BerriAI/litellm.git
synced 2026-09-30 01:52:18 +00:00
perf(vertex/gemini): eliminate O(n²) json.loads on accumulated streaming chunks (#26181)
Replace the per-chunk `accumulated_json += chunk` + full json.loads with a list-based buffer plus a cheap brace/bracket completeness heuristic so the buffer is only joined and parsed when the latest shard looks like a closed envelope. The empty-string flush from __next__/__anext__ at StopIteration still forces a final attempt so partial-completion semantics are unchanged. Re-PR of #26187 against the new litellm_internal_staging base (the original PR auto-closed when the daily-dated litellm_oss_staging_04_25_2026 base branch was retired). Scope intentionally narrowed to the Gemini hot path; the parallel fixes for Anthropic chat handler and SageMaker common_utils will follow once this lands so each can be reviewed independently. Signed-off-by: Tai An <antai12232931@outlook.com>
This commit is contained in:
parent
68852ef165
commit
9892459b13
1 changed files with 21 additions and 10 deletions
|
|
@ -3263,7 +3263,12 @@ class ModelResponseIterator:
|
|||
|
||||
self.streaming_response = streaming_response
|
||||
self.chunk_type: Literal["valid_json", "accumulated_json"] = "valid_json"
|
||||
self.accumulated_json = ""
|
||||
# Buffer SSE shards as a list to avoid O(n²) string concatenation when
|
||||
# Gemini emits multi-MB tool-call payloads as many small chunks. The
|
||||
# parser only joins+json.loads when the trailing brace/bracket suggests
|
||||
# a complete envelope, instead of re-parsing the growing buffer on every
|
||||
# chunk and freezing the event loop.
|
||||
self.accumulated_json_chunks: list = []
|
||||
self.sent_first_chunk = False
|
||||
self.logging_obj = logging_obj
|
||||
self.response_headers = response_headers or {}
|
||||
|
|
@ -3486,16 +3491,22 @@ class ModelResponseIterator:
|
|||
chunk = litellm.CustomStreamWrapper._strip_sse_data_from_chunk(chunk) or ""
|
||||
message = chunk.replace("\n\n", "")
|
||||
|
||||
# Accumulate JSON data
|
||||
self.accumulated_json += message
|
||||
|
||||
# Try to parse the accumulated JSON
|
||||
self.accumulated_json_chunks.append(message)
|
||||
# Cheap completeness heuristic: only join + json.loads when the latest
|
||||
# shard looks like it might have closed the envelope. Avoids parsing
|
||||
# the growing buffer on every chunk -- the empty-string flush from
|
||||
# __next__/__anext__ at StopIteration still forces a final attempt.
|
||||
_stripped = message.rstrip()
|
||||
if _stripped and _stripped[-1] not in ('}', ']'):
|
||||
return None
|
||||
_full_json = "".join(self.accumulated_json_chunks)
|
||||
if not _full_json:
|
||||
return None
|
||||
try:
|
||||
_data = json.loads(self.accumulated_json)
|
||||
self.accumulated_json = "" # reset after successful parsing
|
||||
_data = json.loads(_full_json)
|
||||
self.accumulated_json_chunks = []
|
||||
return self.chunk_parser(chunk=_data)
|
||||
except json.JSONDecodeError:
|
||||
# If it's not valid JSON yet, continue to the next event
|
||||
return None
|
||||
|
||||
def _common_chunk_parsing_logic(
|
||||
|
|
@ -3522,7 +3533,7 @@ class ModelResponseIterator:
|
|||
try:
|
||||
chunk = self.response_iterator.__next__()
|
||||
except StopIteration:
|
||||
if self.chunk_type == "accumulated_json" and self.accumulated_json:
|
||||
if self.chunk_type == "accumulated_json" and self.accumulated_json_chunks:
|
||||
return self.handle_accumulated_json_chunk(chunk="")
|
||||
raise StopIteration
|
||||
except ValueError as e:
|
||||
|
|
@ -3544,7 +3555,7 @@ class ModelResponseIterator:
|
|||
try:
|
||||
chunk = await self.async_response_iterator.__anext__()
|
||||
except StopAsyncIteration:
|
||||
if self.chunk_type == "accumulated_json" and self.accumulated_json:
|
||||
if self.chunk_type == "accumulated_json" and self.accumulated_json_chunks:
|
||||
return self.handle_accumulated_json_chunk(chunk="")
|
||||
raise StopAsyncIteration
|
||||
except ValueError as e:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue