fix(vertex_ai/files): stream OpenAI->Vertex batch JSONL transform to fix OOM on large Gemini batch uploads

This commit is contained in:
Yassin Kortam 2026-05-29 12:56:38 -07:00
parent 68852ef165
commit c2a0d48036
3 changed files with 455 additions and 148 deletions

View file

@ -3,7 +3,7 @@ import json
import os
import re
import time
from typing import Any, Callable, Dict, List, Optional, Tuple, Union
from typing import Any, Callable, Dict, Iterator, List, Optional, Tuple, Union, cast
import httpx
from httpx import Headers, Response
@ -137,42 +137,150 @@ def _get_litellm_batch_custom_id_from_labels(labels: Dict[str, Any]) -> str:
return str(labels.get("litellm_custom_id", "unknown"))
def _openai_batch_jsonl_entries_to_vertex_wrapped_requests(
openai_jsonl_content: List[Dict[str, Any]],
def _openai_batch_jsonl_entry_to_vertex_wrapped_request(
openai_entry: Dict[str, Any],
map_openai_to_vertex_params: Callable[[Dict[str, Any]], Dict[str, Any]],
) -> List[Dict[str, Any]]:
) -> Dict[str, Any]:
"""
Transforms OpenAI JSONL batch entries to Vertex AI JSONL lines.
Transforms a single OpenAI JSONL batch entry into its Vertex wrapped request.
jsonl body for vertex is {"request": <request_body>}
Example Vertex jsonl
{"request":{"contents": [{"role": "user", "parts": [{"text": "What is the relation between the following video and image samples?"}, {"fileData": {"fileUri": "gs://cloud-samples-data/generative-ai/video/animals.mp4", "mimeType": "video/mp4"}}, {"fileData": {"fileUri": "gs://cloud-samples-data/generative-ai/image/cricket.jpeg", "mimeType": "image/jpeg"}}]}]}}
{"request":{"contents": [{"role": "user", "parts": [{"text": "Describe what is happening in this video."}, {"fileData": {"fileUri": "gs://cloud-samples-data/generative-ai/video/another_video.mov", "mimeType": "video/mov"}}]}]}}
"""
openai_request_body = openai_entry.get("body") or {}
vertex_request_body = _transform_request_body(
messages=openai_request_body.get("messages", []),
model=openai_request_body.get("model", ""),
optional_params=map_openai_to_vertex_params(openai_request_body),
custom_llm_provider="vertex_ai",
litellm_params={},
cached_content=None,
)
vertex_jsonl_content = []
for _openai_jsonl_content in openai_jsonl_content:
openai_request_body = _openai_jsonl_content.get("body") or {}
vertex_request_body = _transform_request_body(
messages=openai_request_body.get("messages", []),
model=openai_request_body.get("model", ""),
optional_params=map_openai_to_vertex_params(openai_request_body),
custom_llm_provider="vertex_ai",
litellm_params={},
cached_content=None,
custom_id = openai_entry.get("custom_id")
if custom_id is not None:
if "labels" not in vertex_request_body:
vertex_request_body["labels"] = {}
_set_litellm_batch_custom_id_labels(vertex_request_body["labels"], custom_id)
return {"request": vertex_request_body}
def _openai_batch_jsonl_entries_to_vertex_wrapped_requests(
openai_jsonl_content: List[Dict[str, Any]],
map_openai_to_vertex_params: Callable[[Dict[str, Any]], Dict[str, Any]],
) -> List[Dict[str, Any]]:
return [
_openai_batch_jsonl_entry_to_vertex_wrapped_request(
entry, map_openai_to_vertex_params
)
for entry in openai_jsonl_content
]
# Add custom_id as a label for correlation in batch outputs
custom_id = _openai_jsonl_content.get("custom_id")
if custom_id is not None:
if "labels" not in vertex_request_body:
vertex_request_body["labels"] = {}
_set_litellm_batch_custom_id_labels(
vertex_request_body["labels"], custom_id
def _iter_openai_jsonl_lines(openai_file_content: FileTypes) -> Iterator[str]:
"""
Yield non-empty JSONL lines one at a time without materializing the whole
payload, so peak memory stays bounded regardless of payload size. Mirrors
``str.splitlines()`` + ``line.strip()`` for ``\\n`` / ``\\r\\n`` delimited
JSONL.
"""
content: Any = openai_file_content
if isinstance(content, tuple):
content = content[1]
if isinstance(content, PathLike):
with open(str(content), "rb") as handle:
for raw in handle:
line = raw.decode("utf-8").strip()
if line:
yield line
return
if isinstance(content, (bytes, bytearray)):
newline = ord("\n")
start, length = 0, len(content)
while start < length:
idx = content.find(newline, start)
if idx == -1:
chunk, start = content[start:], length
else:
chunk, start = content[start:idx], idx + 1
line = chunk.decode("utf-8").strip()
if line:
yield line
return
if isinstance(content, str):
str_start, str_length = 0, len(content)
while str_start < str_length:
str_idx = content.find("\n", str_start)
if str_idx == -1:
str_chunk, str_start = content[str_start:], str_length
else:
str_chunk, str_start = content[str_start:str_idx], str_idx + 1
str_line = str_chunk.strip()
if str_line:
yield str_line
return
if hasattr(content, "read"):
for raw in content:
line = raw.decode("utf-8") if isinstance(raw, bytes) else raw
line = line.strip()
if line:
yield line
return
raise ValueError("Unsupported file content type")
def _iter_openai_jsonl_entries(
openai_file_content: FileTypes,
) -> Iterator[Dict[str, Any]]:
for line in _iter_openai_jsonl_lines(openai_file_content):
yield json.loads(line)
def _stream_openai_jsonl_to_vertex(
openai_file_content: FileTypes,
map_openai_to_vertex_params: Callable[[Dict[str, Any]], Dict[str, Any]],
as_bytes: bool,
) -> Tuple[Union[str, bytes], Optional[Dict[str, Any]]]:
"""
Stream the OpenAI -> Vertex JSONL transform entry-by-entry.
Returns the joined Vertex JSONL payload and the first parsed OpenAI entry,
so callers can derive the GCS object name without re-parsing. ``as_bytes``
returns bytes (the upload path ships them to httpx without a str->bytes
re-encode); otherwise a str is returned.
The bytes path extends one ``bytearray`` in place instead of collecting a
list of encoded lines and joining, which would hold both the list and the
joined result at once and double peak output memory.
"""
first_entry: Optional[Dict[str, Any]] = None
byte_buf = bytearray()
str_parts: List[str] = []
for entry in _iter_openai_jsonl_entries(openai_file_content):
if first_entry is None:
first_entry = entry
line = json.dumps(
_openai_batch_jsonl_entry_to_vertex_wrapped_request(
entry, map_openai_to_vertex_params
)
)
if as_bytes:
if byte_buf:
byte_buf.extend(b"\n")
byte_buf.extend(line.encode("utf-8"))
else:
str_parts.append(line)
vertex_jsonl_content.append({"request": vertex_request_body})
return vertex_jsonl_content
if as_bytes:
return bytes(byte_buf), first_entry
return "\n".join(str_parts), first_entry
class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
@ -208,43 +316,6 @@ class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
headers["Authorization"] = f"Bearer {api_key}"
return headers
def _get_content_from_openai_file(self, openai_file_content: FileTypes) -> str:
"""
Helper to extract content from various OpenAI file types and return as string.
Handles:
- Direct content (str, bytes, IO[bytes])
- Tuple formats: (filename, content, [content_type], [headers])
- PathLike objects
"""
content: Union[str, bytes] = b""
# Extract file content from tuple if necessary
if isinstance(openai_file_content, tuple):
# Take the second element which is always the file content
file_content = openai_file_content[1]
else:
file_content = openai_file_content
# Handle different file content types
if isinstance(file_content, str):
# String content can be used directly
content = file_content
elif isinstance(file_content, bytes):
# Bytes content can be decoded
content = file_content
elif isinstance(file_content, PathLike): # PathLike
with open(str(file_content), "rb") as f:
content = f.read()
elif hasattr(file_content, "read"): # IO[bytes]
# File-like objects need to be read
content = file_content.read()
# Ensure content is string
if isinstance(content, bytes):
content = content.decode("utf-8")
return content
def _get_gcs_object_name_from_batch_jsonl(
self,
openai_jsonl_content: List[Dict[str, Any]],
@ -273,17 +344,12 @@ class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
raise ValueError("file content is required")
if purpose == "batch":
## 1. If jsonl, check if there's a model name
file_content = self._get_content_from_openai_file(
extracted_file_data_content
## 1. If jsonl, derive the object name from the first entry's model
first_entry = next(
_iter_openai_jsonl_entries(extracted_file_data_content), None
)
# Split into lines and parse each line as JSON
openai_jsonl_content = [
json.loads(line) for line in file_content.splitlines() if line.strip()
]
if len(openai_jsonl_content) > 0:
return self._get_gcs_object_name_from_batch_jsonl(openai_jsonl_content)
if first_entry is not None:
return self._get_gcs_object_name_from_batch_jsonl([first_entry])
## 2. If not jsonl, store under a server-generated managed object name
filename = extracted_file_data.get("filename")
@ -399,21 +465,12 @@ class VertexAIFilesConfig(VertexBase, BaseFilesConfig):
create_file_data=create_file_data,
extracted_file_data=extracted_file_data,
):
## 1. If jsonl, check if there's a model name
file_content = self._get_content_from_openai_file(
extracted_file_data_content
vertex_jsonl_bytes, _ = _stream_openai_jsonl_to_vertex(
extracted_file_data_content,
self._map_openai_to_vertex_params,
as_bytes=True,
)
# Split into lines and parse each line as JSON
openai_jsonl_content = [
json.loads(line) for line in file_content.splitlines() if line.strip()
]
vertex_jsonl_content = (
self._transform_openai_jsonl_content_to_vertex_ai_jsonl_content(
openai_jsonl_content
)
)
return "\n".join(json.dumps(item) for item in vertex_jsonl_content)
return vertex_jsonl_bytes
elif isinstance(extracted_file_data_content, bytes):
return extracted_file_data_content
else:
@ -811,25 +868,16 @@ class VertexAIJsonlFilesTransformation(VertexGeminiConfig):
if openai_file_content is None:
raise ValueError("contents of file are None")
# Read the content of the file
file_content = self._get_content_from_openai_file(openai_file_content)
# Split into lines and parse each line as JSON
openai_jsonl_content = [
json.loads(line) for line in file_content.splitlines() if line.strip()
]
vertex_jsonl_content = (
self._transform_openai_jsonl_content_to_vertex_ai_jsonl_content(
openai_jsonl_content
)
vertex_jsonl_string, first_entry = _stream_openai_jsonl_to_vertex(
openai_file_content,
self._map_openai_to_vertex_params,
as_bytes=False,
)
vertex_jsonl_string = "\n".join(
json.dumps(item) for item in vertex_jsonl_content
)
object_name = self._get_gcs_object_name(
openai_jsonl_content=openai_jsonl_content
)
return vertex_jsonl_string, object_name
if first_entry is None:
raise ValueError("contents of file are empty")
object_name = self._get_gcs_object_name(openai_jsonl_content=[first_entry])
return cast(str, vertex_jsonl_string), object_name
def _transform_openai_jsonl_content_to_vertex_ai_jsonl_content(
self, openai_jsonl_content: List[Dict[str, Any]]
@ -871,43 +919,6 @@ class VertexAIJsonlFilesTransformation(VertexGeminiConfig):
)
return vertex_params
def _get_content_from_openai_file(self, openai_file_content: FileTypes) -> str:
"""
Helper to extract content from various OpenAI file types and return as string.
Handles:
- Direct content (str, bytes, IO[bytes])
- Tuple formats: (filename, content, [content_type], [headers])
- PathLike objects
"""
content: Union[str, bytes] = b""
# Extract file content from tuple if necessary
if isinstance(openai_file_content, tuple):
# Take the second element which is always the file content
file_content = openai_file_content[1]
else:
file_content = openai_file_content
# Handle different file content types
if isinstance(file_content, str):
# String content can be used directly
content = file_content
elif isinstance(file_content, bytes):
# Bytes content can be decoded
content = file_content
elif isinstance(file_content, PathLike): # PathLike
with open(str(file_content), "rb") as f:
content = f.read()
elif hasattr(file_content, "read"): # IO[bytes]
# File-like objects need to be read
content = file_content.read()
# Ensure content is string
if isinstance(content, bytes):
content = content.decode("utf-8")
return content
def transform_gcs_bucket_response_to_openai_file_object(
self, create_file_data: CreateFileRequest, gcs_upload_response: Dict[str, Any]
) -> OpenAIFileObject:

View file

@ -8,6 +8,7 @@ Regression test for: UTF-8 codec error when uploading binary files
"""
import io
import json
import pytest
from unittest.mock import AsyncMock, MagicMock, patch
@ -137,11 +138,12 @@ class TestVertexAIBinaryFileUpload:
), "Binary file data should remain as bytes"
@pytest.mark.asyncio
async def test_jsonl_file_upload_returns_string(self):
async def test_jsonl_file_upload_returns_utf8_bytes(self):
"""
Test that JSONL files (text) are correctly transformed to strings.
Test that JSONL batch files are transformed to UTF-8 bytes.
This ensures we handle both binary and text files correctly.
The transform emits bytes so the upload ships the payload to GCS without
a str->bytes re-encode, avoiding a full extra copy of the payload.
"""
# Create mock JSONL content
mock_jsonl_content = (
@ -164,10 +166,14 @@ class TestVertexAIBinaryFileUpload:
litellm_params={},
)
# JSONL files should be transformed to string
assert isinstance(
transformed_request, str
), f"Expected string for JSONL file, got {type(transformed_request)}"
transformed_request, bytes
), f"Expected bytes for JSONL file, got {type(transformed_request)}"
decoded = json.loads(transformed_request.decode("utf-8"))
assert (
"request" in decoded
), "JSONL transform must wrap each row in {'request': ...}"
@pytest.mark.asyncio
async def test_mixed_file_types_in_sequence(self):
@ -208,7 +214,7 @@ class TestVertexAIBinaryFileUpload:
optional_params={},
litellm_params={},
)
assert isinstance(result2, str)
assert isinstance(result2, bytes)
# Test 3: Upload another binary file
binary_content2 = b"\xc4\xe5\xf2\xe5\xeb"
@ -251,7 +257,7 @@ class TestVertexAIBinaryFileUpload:
},
"text_files": {
"input_type": "str or bytes",
"output_type": "str",
"output_type": "bytes",
"examples": ["JSONL", "CSV", "TXT"],
"http_method": "POST",
"encoding": "UTF-8",

View file

@ -0,0 +1,290 @@
"""
Tests for the streaming OpenAI -> Vertex JSONL batch transform.
The transform converts batch uploads entry-by-entry rather than materializing
the payload in full intermediate lists (decoded str, parsed dicts, transformed
dicts, joined output), which keeps peak memory bounded on large uploads.
These tests lock in the behaviour that would regress if the streaming path were
replaced by a list-based pipeline:
1. Byte-for-byte output parity with a list pipeline (wire format).
2. The streaming transform peaks at a clear fraction of a list pipeline on the
same input (relative differential, robust to GC noise).
3. ``get_object_name`` only parses the first JSONL row, so a payload whose
later rows are not valid JSON does not raise.
4. A tuple-wrapped file handle uploaded through the real create_file ordering
keeps every row, including entry 0 (no partial upload from a consumed
cursor).
"""
import gc
import io
import json
import tracemalloc
import pytest
from litellm.llms.vertex_ai.files.transformation import (
VertexAIFilesConfig,
VertexAIJsonlFilesTransformation,
_get_litellm_batch_custom_id_from_labels,
_iter_openai_jsonl_entries,
_iter_openai_jsonl_lines,
_stream_openai_jsonl_to_vertex,
)
from litellm.types.llms.openai import CreateFileRequest
def _make_openai_jsonl_bytes(n_rows: int, padding: int = 400) -> bytes:
pad = "x" * padding
rows = []
for i in range(n_rows):
rows.append(
json.dumps(
{
"custom_id": f"request-{i}",
"method": "POST",
"url": "/v1/chat/completions",
"body": {
"model": "gemini-2.5-flash",
"messages": [{"role": "user", "content": f"{pad} {i}"}],
"max_tokens": 4,
},
}
)
)
return ("\n".join(rows)).encode("utf-8")
def _legacy_vertex_jsonl_string(cfg: VertexAIFilesConfig, content: str) -> str:
"""A list-based pipeline, reconstructed for parity comparison."""
entries = [json.loads(line) for line in content.splitlines() if line.strip()]
vertex = cfg._transform_openai_jsonl_content_to_vertex_ai_jsonl_content(entries)
return "\n".join(json.dumps(item) for item in vertex)
class TestStreamingOutputParity:
def test_streaming_bytes_match_legacy_pipeline(self):
cfg = VertexAIFilesConfig()
raw = _make_openai_jsonl_bytes(500)
streamed, first_entry = _stream_openai_jsonl_to_vertex(
raw, cfg._map_openai_to_vertex_params, as_bytes=True
)
legacy = _legacy_vertex_jsonl_string(cfg, raw.decode("utf-8"))
assert isinstance(streamed, bytes)
assert streamed.decode("utf-8") == legacy
assert first_entry is not None
assert first_entry["custom_id"] == "request-0"
def test_transform_create_file_request_returns_bytes_parity(self):
cfg = VertexAIFilesConfig()
raw = _make_openai_jsonl_bytes(300)
request: CreateFileRequest = {
"file": ("batch.jsonl", raw, "application/jsonl"),
"purpose": "batch",
}
out = cfg.transform_create_file_request(
model="", create_file_data=request, optional_params={}, litellm_params={}
)
assert isinstance(out, bytes)
assert out.decode("utf-8") == _legacy_vertex_jsonl_string(
cfg, raw.decode("utf-8")
)
class TestFileLikeInputNotPartiallyConsumed:
"""
In ``llm_http_handler.create_file`` the object-name step
(get_complete_file_url -> get_object_name) runs before
transform_create_file_request, and each independently calls
``extract_file_data`` on the same create_file_data. When the file is a
tuple-wrapped open handle, the streaming reader must still emit every row
including entry 0: ``extract_file_data`` materializes the handle to bytes and
rewinds it (seek(0)), so neither step consumes the other's cursor. A partial
upload missing the first request would be silent and hard to catch, so this
locks the full-payload invariant in.
"""
def test_filehandle_create_file_keeps_first_entry(self):
cfg = VertexAIFilesConfig()
n_rows = 25
raw = _make_openai_jsonl_bytes(n_rows)
create_file_data: CreateFileRequest = {
"file": ("batch.jsonl", io.BytesIO(raw), "application/jsonl"),
"purpose": "batch",
}
# Object-name step first (as the handler does), then the transform, both
# reading the same live BytesIO handle.
cfg.get_complete_file_url(
api_base=None,
api_key=None,
model="",
optional_params={},
litellm_params={"bucket_name": "test-bucket"},
data=create_file_data,
)
out = cfg.transform_create_file_request(
model="",
create_file_data=create_file_data,
optional_params={},
litellm_params={},
)
assert isinstance(out, bytes)
lines = out.decode("utf-8").splitlines()
assert len(lines) == n_rows, "no batch row may be dropped from the upload"
first_labels = json.loads(lines[0])["request"]["labels"]
assert _get_litellm_batch_custom_id_from_labels(first_labels) == "request-0"
class TestStreamingLineIterator:
def test_skips_blank_and_whitespace_lines(self):
content = b'{"a": 1}\n\n \n{"b": 2}\n'
assert list(_iter_openai_jsonl_lines(content)) == ['{"a": 1}', '{"b": 2}']
def test_handles_crlf_and_missing_trailing_newline(self):
content = b'{"a": 1}\r\n{"b": 2}'
assert [json.loads(line) for line in _iter_openai_jsonl_lines(content)] == [
{"a": 1},
{"b": 2},
]
def test_accepts_str_bytes_tuple_and_filelike(self):
expected = [{"a": 1}, {"b": 2}]
text = '{"a": 1}\n{"b": 2}\n'
for source in (
text,
text.encode("utf-8"),
("name.jsonl", text.encode("utf-8"), "application/jsonl"),
io.BytesIO(text.encode("utf-8")),
):
assert list(_iter_openai_jsonl_entries(source)) == expected
def test_str_input_without_trailing_newline(self):
assert list(_iter_openai_jsonl_lines('{"a": 1}\n{"b": 2}')) == [
'{"a": 1}',
'{"b": 2}',
]
def test_pathlike_input_is_read_line_by_line(self, tmp_path):
path = tmp_path / "batch.jsonl"
path.write_bytes(b'{"a": 1}\n{"b": 2}\n')
assert list(_iter_openai_jsonl_entries(path)) == [{"a": 1}, {"b": 2}]
def test_unsupported_content_type_raises(self):
with pytest.raises(ValueError, match="Unsupported file content type"):
list(_iter_openai_jsonl_lines(12345)) # type: ignore[arg-type]
def test_is_lazy_does_not_parse_past_first_entry(self):
# Second row is invalid JSON; pulling only the first entry must not raise.
content = b'{"custom_id": "first"}\nnot-json-at-all\n'
gen = _iter_openai_jsonl_entries(content)
assert next(gen)["custom_id"] == "first"
with pytest.raises(json.JSONDecodeError):
next(gen)
class TestLegacyHandlerPathStreaming:
def test_returns_str_with_object_name_from_first_row(self):
transformer = VertexAIJsonlFilesTransformation()
raw = _make_openai_jsonl_bytes(50)
vertex_str, object_name = (
transformer.transform_openai_file_content_to_vertex_ai_file_content(raw)
)
assert isinstance(vertex_str, str)
assert "gemini-2.5-flash" in object_name
# First line is a valid Vertex-wrapped request.
assert "request" in json.loads(vertex_str.splitlines()[0])
def test_empty_payload_raises(self):
transformer = VertexAIJsonlFilesTransformation()
with pytest.raises(ValueError, match="empty"):
transformer.transform_openai_file_content_to_vertex_ai_file_content(b"\n\n")
class TestGetObjectNameLazyParse:
def test_only_parses_first_row_for_model(self):
cfg = VertexAIFilesConfig()
# Tail rows are deliberately not valid JSON. Parsing the whole payload
# would raise here; a first-row-only parse must not.
raw = (
b'{"custom_id": "r-0", "body": {"model": "gemini-2.5-flash"}}\n'
b"garbage line that is not json\n"
)
from litellm.litellm_core_utils.prompt_templates.common_utils import (
extract_file_data,
)
extracted = extract_file_data(("batch.jsonl", raw, "application/jsonl"))
object_name = cfg.get_object_name(extracted, purpose="batch")
assert "gemini-2.5-flash" in object_name
class TestStreamingPeakMemory:
"""
Differential guard: the streaming transform must stay well under the peak
that a list pipeline incurs on the same input. If the hot path builds full
intermediate lists, the streaming assertion fails.
The assertion that matters is the *relative* one: ``streaming_peak`` must be
a clear fraction of ``list_peak`` on the identical input. Absolute
``tracemalloc`` ratios drift with GC timing and the live set carried in from
earlier tests, so they make poor CI gates; the relative comparison cancels
that shared noise and is exactly what regresses (toward 1.0) when the hot
path builds full intermediate lists. ``gc.collect()`` before each
measurement removes any garbage the previous run left behind.
"""
def _measure(self, fn):
gc.collect()
tracemalloc.start()
try:
fn()
_, peak = tracemalloc.get_traced_memory()
finally:
tracemalloc.stop()
return peak
def test_streaming_peak_well_below_list_pipeline(self):
cfg = VertexAIFilesConfig()
raw = _make_openai_jsonl_bytes(8000)
content_str = raw.decode("utf-8")
streaming_peak = self._measure(
lambda: _stream_openai_jsonl_to_vertex(
raw, cfg._map_openai_to_vertex_params, as_bytes=True
)
)
list_peak = self._measure(lambda: _legacy_vertex_jsonl_string(cfg, content_str))
# Core guard: streaming peaks at well under two-thirds of the list
# pipeline. Building full intermediate lists in the hot path pushes this
# ratio back toward 1.0 and fails the test.
assert streaming_peak < list_peak * 0.6, (
f"streaming peak {streaming_peak} not a clear win over list pipeline "
f"{list_peak} (ratio {streaming_peak / list_peak:.2f})"
)
def test_get_object_name_does_not_scale_with_payload(self):
cfg = VertexAIFilesConfig()
from litellm.litellm_core_utils.prompt_templates.common_utils import (
extract_file_data,
)
raw = _make_openai_jsonl_bytes(8000)
extracted = extract_file_data(("batch.jsonl", raw, "application/jsonl"))
# The payload bytes already exist before measurement starts, so a lazy
# first-row parse should allocate only a small fraction of the payload;
# parsing every row would blow past this bound.
peak = self._measure(lambda: cfg.get_object_name(extracted, purpose="batch"))
assert (
peak / len(raw) < 2.0
), "get_object_name should not copy the whole payload"