From 9892459b13cff2d16e8c03f0cbb34070f34483f7 Mon Sep 17 00:00:00 2001 From: Tai An Date: Fri, 29 May 2026 09:18:10 -0700 Subject: [PATCH] =?UTF-8?q?perf(vertex/gemini):=20eliminate=20O(n=C2=B2)?= =?UTF-8?q?=20json.loads=20on=20accumulated=20streaming=20chunks=20(#26181?= =?UTF-8?q?)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../vertex_and_google_ai_studio_gemini.py | 31 +++++++++++++------ 1 file changed, 21 insertions(+), 10 deletions(-) diff --git a/litellm/llms/vertex_ai/gemini/vertex_and_google_ai_studio_gemini.py b/litellm/llms/vertex_ai/gemini/vertex_and_google_ai_studio_gemini.py index 189ac7a7f6a..03cd30fa48f 100644 --- a/litellm/llms/vertex_ai/gemini/vertex_and_google_ai_studio_gemini.py +++ b/litellm/llms/vertex_ai/gemini/vertex_and_google_ai_studio_gemini.py @@ -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: