mirror of
https://github.com/PaddlePaddle/FastDeploy.git
synced 2025-12-24 13:28:13 +08:00
[Feature] mm support prefix cache (#4134)
* support mm prefix caching * update code * fix mm_hashes * support encoder cache * add encoder cache * update code * update encoder cache * fix features bug * fix worker bug * support processor cache, need to optimize yet * refactor multimodal data cache * update code * update code * update v1 scheduler * update code * update code * update codestyle * support turn off processor cache and encoder cache * update pre-commit * fix code * solve review * update code * update code * update test case * set processor cache in GiB * update test case * support mm prefix caching for qwen model * fix code style check * update pre-commit * fix unit test * fix unit test * add ci test case * fix rescheduled bug * change text_after_process to prompt_tokens * fix unit test * fix chat template * change model path * [EP] fix adapter bugs (#4572) * Update expert_service.py * Update common_engine.py * Update expert_service.py * fix v1 hang bug (#4573) * fix import image_ops error on some platforms (#4559) * [CLI]Update parameters in bench latecy cli tool and fix collect-env cli tool (#4558) * add collect-env * del files * [Graph Optimization] Add dy_runnable and introduce cudagraph_switch_threshold for cudagraph mode switching (#4578) * add new branch for sot * reorder * fix batch bug * [XPU]Moe uses a new operator (#4585) * [XPU]Moe uses a new operator * [XPU]Moe uses a new operator * update response * [Feature] Support Paddle-OCR (#4396) * init * update code * fix code style & disable thinking * adapt for common_engine.update_mm_requests_chunk_size * use 3d rope * use flash_attn_unpadded * opt siglip * update to be compatible with the latest codebase * fix typo * optim OCR performance * fix bug * fix bug * fix bug * fix bug * normlize name * modify xpu rope * revert logger * fix bug * fix bug * fix bug * support default_v1 * optim performance * fix bug --------- Co-authored-by: root <root@szzj-acg-tge1-fdda9.szzj.baidu.com> Co-authored-by: zhangyue66 <zhangyue66@baidu.com> * [DataProcessor] add reasoning_tokens into usage info (#4520) * add reasoning_tokens into usage info initial commit * add unit tests * modify unit test * modify and add unit tests * fix unit test * move steam usage to processor * modify processor * modify test_logprobs * modify test_logprobs.py * modify stream reasoning tokens accumulation * fix unit test * perf: Optimize task queue communication from engine to worker (#4531) * perf: Optimize task queue communication from engine to worker * perf: get_tasks to numpy * perf: get_tasks remove to_numpy * fix: request & replace ENV * remove test_e2w_perf.py * fix code style --------- Co-authored-by: Jiang-Jia-Jun <163579578+Jiang-Jia-Jun@users.noreply.github.com> * Clean up ports after processing results (#4587) * [CI] Add /re-run command in PR comments to restart failed CI workflows (#4593) * [Others] api server exits when worker process is dead (#3271) * [fix] fix terminal hangs when worker process is dead * [chore] change sleep time of monitor * [chore] remove redundant comments * update docs --------- Co-authored-by: ApplEOFDiscord <wwy640130@163.com> Co-authored-by: ApplEOFDiscord <31272106+ApplEOFDiscord@users.noreply.github.com> Co-authored-by: ltd0924 <32387785+ltd0924@users.noreply.github.com> Co-authored-by: yinwei <yinwei_hust@163.com> Co-authored-by: JYChen <zoooo0820@qq.com> Co-authored-by: qwes5s5 <45442318+qwes5s5@users.noreply.github.com> Co-authored-by: Ryan <zihaohuang@aliyun.com> Co-authored-by: yyssys <atyangshuang@foxmail.com> Co-authored-by: ming1753 <61511741+ming1753@users.noreply.github.com> Co-authored-by: root <root@szzj-acg-tge1-fdda9.szzj.baidu.com> Co-authored-by: zhangyue66 <zhangyue66@baidu.com> Co-authored-by: kxz2002 <115912648+kxz2002@users.noreply.github.com> Co-authored-by: SunLei <sunlei5788@gmail.com> Co-authored-by: Jiang-Jia-Jun <163579578+Jiang-Jia-Jun@users.noreply.github.com> Co-authored-by: Zhang Yulong <35552275+ZhangYulongg@users.noreply.github.com> Co-authored-by: YuBaoku <49938469+EmmonsCurse@users.noreply.github.com> Co-authored-by: 李泳桦 <39643373+liyonghua0910@users.noreply.github.com>
This commit is contained in:
@@ -17,7 +17,6 @@
|
||||
import os
|
||||
import time
|
||||
import uuid
|
||||
from copy import deepcopy
|
||||
from pathlib import Path
|
||||
from typing import List, Literal, Optional, Union
|
||||
from urllib.parse import urlparse
|
||||
@@ -29,6 +28,7 @@ from openai.types.chat import (
|
||||
from openai.types.chat import (
|
||||
ChatCompletionMessageParam as OpenAIChatCompletionMessageParam,
|
||||
)
|
||||
from openai.types.chat.chat_completion_content_part_image_param import ImageURL
|
||||
from typing_extensions import Required, TypeAlias, TypedDict
|
||||
|
||||
from fastdeploy.multimodal.image import ImageMediaIO
|
||||
@@ -36,6 +36,17 @@ from fastdeploy.multimodal.video import VideoMediaIO
|
||||
from fastdeploy.utils import api_server_logger
|
||||
|
||||
|
||||
class CustomChatCompletionContentPartImageParam(TypedDict, total=False):
|
||||
"""Custom Image URL object"""
|
||||
|
||||
type: Required[Literal["image_url"]]
|
||||
"""The type of the content part."""
|
||||
|
||||
image_url: Optional[ImageURL]
|
||||
|
||||
uuid: Optional[str]
|
||||
|
||||
|
||||
class VideoURL(TypedDict, total=False):
|
||||
"""Video URL object"""
|
||||
|
||||
@@ -46,14 +57,17 @@ class VideoURL(TypedDict, total=False):
|
||||
class CustomChatCompletionContentPartVideoParam(TypedDict, total=False):
|
||||
"""Custom Video URL object"""
|
||||
|
||||
video_url: Required[VideoURL]
|
||||
|
||||
type: Required[Literal["video_url"]]
|
||||
"""The type of the content type."""
|
||||
"""The type of the content part."""
|
||||
|
||||
video_url: Optional[VideoURL]
|
||||
|
||||
uuid: Optional[str]
|
||||
|
||||
|
||||
CustomChatCompletionContentPartParam: TypeAlias = Union[
|
||||
OpenAIChatCompletionContentPartParam,
|
||||
CustomChatCompletionContentPartImageParam,
|
||||
CustomChatCompletionContentPartVideoParam,
|
||||
]
|
||||
|
||||
@@ -77,7 +91,7 @@ class CustomChatCompletionMessageParam(TypedDict, total=False):
|
||||
ChatCompletionMessageParam = Union[OpenAIChatCompletionMessageParam, CustomChatCompletionMessageParam]
|
||||
|
||||
|
||||
class MultiModalPartParser:
|
||||
class MultimodalPartParser:
|
||||
"""Multi Modal Part parser"""
|
||||
|
||||
def __init__(self):
|
||||
@@ -139,32 +153,46 @@ def parse_content_part(mm_parser, part):
|
||||
return part
|
||||
|
||||
if part_type == "image_url":
|
||||
content = part.get("image_url", {}).get("url", None)
|
||||
image = mm_parser.parse_image(content)
|
||||
parsed = deepcopy(part)
|
||||
del parsed["image_url"]["url"]
|
||||
parsed["image"] = image
|
||||
parsed["type"] = "image"
|
||||
return parsed
|
||||
if not part.get("image_url", None) and not part.get("uuid", None):
|
||||
raise ValueError("Both image_url and uuid are missing")
|
||||
|
||||
if part.get("image_url", None):
|
||||
url = part["image_url"]["url"]
|
||||
image = mm_parser.parse_image(url)
|
||||
else:
|
||||
image = None
|
||||
|
||||
parsed = {}
|
||||
parsed["type"] = "image"
|
||||
parsed["data"] = image
|
||||
parsed["uuid"] = part.get("uuid", None)
|
||||
|
||||
return parsed
|
||||
if part_type == "video_url":
|
||||
content = part.get("video_url", {}).get("url", None)
|
||||
video = mm_parser.parse_video(content)
|
||||
parsed = deepcopy(part)
|
||||
del parsed["video_url"]["url"]
|
||||
parsed["video"] = video
|
||||
if not part.get("video_url", None) and not part.get("uuid", None):
|
||||
raise ValueError("Both video_url and uuid are missing")
|
||||
|
||||
if part.get("video_url", None):
|
||||
url = part["video_url"]["url"]
|
||||
video = mm_parser.parse_video(url)
|
||||
else:
|
||||
video = None
|
||||
|
||||
parsed = {}
|
||||
parsed["type"] = "video"
|
||||
parsed["data"] = video
|
||||
parsed["uuid"] = part.get("uuid", None)
|
||||
|
||||
return parsed
|
||||
|
||||
raise ValueError(f"Unknown content part type: {part_type}")
|
||||
|
||||
|
||||
# TODO async
|
||||
# def parse_chat_messages(messages: List[ChatCompletionMessageParam]):
|
||||
def parse_chat_messages(messages):
|
||||
def parse_chat_messages(messages: List[ChatCompletionMessageParam]):
|
||||
"""Parse chat messages to [dict]"""
|
||||
|
||||
mm_parser = MultiModalPartParser()
|
||||
mm_parser = MultimodalPartParser()
|
||||
|
||||
conversation = []
|
||||
for message in messages:
|
||||
|
||||
@@ -68,16 +68,19 @@ class EngineClient:
|
||||
tool_parser=None,
|
||||
enable_prefix_caching=None,
|
||||
splitwise_role=None,
|
||||
max_processor_cache=0,
|
||||
):
|
||||
model_config = ModelConfig({"model": model_name_or_path})
|
||||
self.enable_mm = model_config.enable_mm
|
||||
enable_processor_cache = self.enable_mm and max_processor_cache > 0
|
||||
input_processor = InputPreprocessor(
|
||||
model_config,
|
||||
reasoning_parser,
|
||||
limit_mm_per_prompt,
|
||||
mm_processor_kwargs,
|
||||
tool_parser,
|
||||
enable_processor_cache,
|
||||
)
|
||||
self.enable_mm = model_config.enable_mm
|
||||
self.enable_logprob = enable_logprob
|
||||
self.reasoning_parser = reasoning_parser
|
||||
self.data_processor = input_processor.create_processor()
|
||||
|
||||
@@ -193,6 +193,7 @@ async def lifespan(app: FastAPI):
|
||||
tool_parser=args.tool_call_parser,
|
||||
enable_prefix_caching=args.enable_prefix_caching,
|
||||
splitwise_role=args.splitwise_role,
|
||||
max_processor_cache=args.max_processor_cache,
|
||||
)
|
||||
await engine_client.connection_manager.initialize()
|
||||
app.state.dynamic_load_weight = args.dynamic_load_weight
|
||||
|
||||
@@ -470,6 +470,8 @@ class CompletionRequest(BaseModel):
|
||||
max_streaming_response_tokens: Optional[int] = None
|
||||
return_token_ids: Optional[bool] = None
|
||||
prompt_token_ids: Optional[Union[List[int], List[List[int]]]] = None
|
||||
|
||||
mm_hashes: Optional[list] = None
|
||||
# doc: end-completion-extra-params
|
||||
|
||||
def to_dict_for_infer(self, request_id=None, prompt=None):
|
||||
@@ -527,6 +529,9 @@ class CompletionRequest(BaseModel):
|
||||
if item is not None:
|
||||
req_dict[key] = item
|
||||
|
||||
if self.mm_hashes is not None and len(self.mm_hashes) > 0:
|
||||
req_dict["mm_hashes"] = self.mm_hashes
|
||||
|
||||
return req_dict
|
||||
|
||||
@model_validator(mode="before")
|
||||
@@ -553,6 +558,9 @@ class CompletionRequest(BaseModel):
|
||||
"('guided_json', 'guided_regex', 'guided_choice', 'guided_grammar')."
|
||||
)
|
||||
|
||||
if data.get("mm_hashes", None):
|
||||
assert isinstance(data["mm_hashes"], list), "`mm_hashes` must be a list."
|
||||
|
||||
return data
|
||||
|
||||
|
||||
@@ -618,6 +626,8 @@ class ChatCompletionRequest(BaseModel):
|
||||
prompt_token_ids: Optional[List[int]] = None
|
||||
max_streaming_response_tokens: Optional[int] = None
|
||||
disable_chat_template: Optional[bool] = False
|
||||
|
||||
mm_hashes: Optional[list] = None
|
||||
completion_token_ids: Optional[List[int]] = None
|
||||
# doc: end-chat-completion-extra-params
|
||||
|
||||
@@ -694,6 +704,9 @@ class ChatCompletionRequest(BaseModel):
|
||||
if item is not None:
|
||||
req_dict[key] = item
|
||||
|
||||
if self.mm_hashes is not None and len(self.mm_hashes) > 0:
|
||||
req_dict["mm_hashes"] = self.mm_hashes
|
||||
|
||||
return req_dict
|
||||
|
||||
@model_validator(mode="before")
|
||||
@@ -721,6 +734,9 @@ class ChatCompletionRequest(BaseModel):
|
||||
"('guided_json', 'guided_regex', 'guided_choice', 'guided_grammar', 'structural_tag')."
|
||||
)
|
||||
|
||||
if data.get("mm_hashes", None):
|
||||
assert isinstance(data["mm_hashes"], list), "`mm_hashes` must be a list."
|
||||
|
||||
return data
|
||||
|
||||
@model_validator(mode="before")
|
||||
|
||||
Reference in New Issue
Block a user