mirror of
https://github.com/BerriAI/litellm.git
synced 2026-10-09 03:18:44 +00:00
fix: async two-step upload uses thread-executor; read full JSONL first line
This commit is contained in:
parent
6f11163038
commit
6a9b84e788
2 changed files with 33 additions and 17 deletions
|
|
@ -152,6 +152,21 @@ else:
|
|||
LiteLLMLoggingObj = Any
|
||||
|
||||
|
||||
async def _async_file_chunks(
|
||||
fp: Any, chunk_size: int = 65536
|
||||
) -> "AsyncIterator[bytes]":
|
||||
"""Async generator that reads a sync file-like object in chunks via a thread executor.
|
||||
|
||||
Dispatches each blocking disk read to a thread pool via anyio so the async
|
||||
event loop is never stalled waiting for I/O.
|
||||
"""
|
||||
while True:
|
||||
chunk = await anyio_to_thread.run_sync(lambda: fp.read(chunk_size))
|
||||
if not chunk:
|
||||
break
|
||||
yield chunk
|
||||
|
||||
|
||||
class BaseLLMHTTPHandler:
|
||||
async def _make_common_async_call(
|
||||
self,
|
||||
|
|
@ -3159,7 +3174,9 @@ class BaseLLMHTTPHandler:
|
|||
"timeout": timeout,
|
||||
}
|
||||
if hasattr(upload_data, "read") and hasattr(upload_data, "seek"):
|
||||
async_upload_kwargs["content"] = upload_data
|
||||
# Wrap sync IO in the async generator so disk reads are
|
||||
# dispatched to a thread executor and do not block the loop.
|
||||
async_upload_kwargs["content"] = _async_file_chunks(upload_data)
|
||||
else:
|
||||
async_upload_kwargs["data"] = upload_data
|
||||
upload_response = await getattr(async_httpx_client, upload_method)(
|
||||
|
|
@ -3202,23 +3219,10 @@ class BaseLLMHTTPHandler:
|
|||
):
|
||||
# Wrap sync IO in an async generator so disk reads are dispatched
|
||||
# to a thread executor and do not block the event loop.
|
||||
file_obj = transformed_request
|
||||
|
||||
async def _async_file_chunks(
|
||||
fp: Any, chunk_size: int = 65536
|
||||
) -> "AsyncIterator[bytes]":
|
||||
while True:
|
||||
chunk = await anyio_to_thread.run_sync(
|
||||
lambda: fp.read(chunk_size)
|
||||
)
|
||||
if not chunk:
|
||||
break
|
||||
yield chunk
|
||||
|
||||
async_upload_kwargs: Dict[str, Any] = {
|
||||
"url": api_base,
|
||||
"headers": headers,
|
||||
"content": _async_file_chunks(file_obj),
|
||||
"content": _async_file_chunks(transformed_request),
|
||||
"timeout": timeout,
|
||||
}
|
||||
else:
|
||||
|
|
|
|||
|
|
@ -452,8 +452,20 @@ async def create_file( # noqa: PLR0915
|
|||
router_model: Optional[str] = None
|
||||
is_router_model = False
|
||||
if litellm.enable_loadbalancing_on_batch_endpoints is True:
|
||||
# Read only the first line to detect the model; seek back afterwards.
|
||||
first_line_bytes = file.file.read(4096)
|
||||
# Read the first complete JSONL line for model detection.
|
||||
# Read in 4096-byte chunks until we find a newline, capping at 1 MB
|
||||
# to avoid loading the entire file for pathologically long lines.
|
||||
_FIRST_LINE_MAX = 1024 * 1024 # 1 MB
|
||||
first_line_bytes = b""
|
||||
while len(first_line_bytes) < _FIRST_LINE_MAX:
|
||||
chunk = file.file.read(4096)
|
||||
if not chunk:
|
||||
break
|
||||
newline_pos = chunk.find(b"\n")
|
||||
if newline_pos != -1:
|
||||
first_line_bytes += chunk[:newline_pos]
|
||||
break
|
||||
first_line_bytes += chunk
|
||||
file.file.seek(0)
|
||||
json_obj = get_first_json_object(file_content_bytes=first_line_bytes)
|
||||
if json_obj:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue