Slackbot
08/09/2023, 10:10 AMJiang
08/09/2023, 10:16 AMJiang
08/09/2023, 10:20 AMJiang
08/09/2023, 10:23 AM@svc.api(input=bentoml.io.Text(), output=bentoml.io.JSON())
async def classify(text: str) -> str:
results = await model_runner.async_run([text])
return results[0]hao chen
08/09/2023, 10:32 AM@bentoml.Runnable.method()
def stream_generate_1(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)
Because the runner method does not support concurrency, all subsequent requests must wait for the current processing to end. I want this runner method to support multiple concurrency so that multiple requests can be processed simultaneously.
My application is the inference service for the large language model, and I am working on Runner's__ Init__ The method has loaded a model that supports concurrent inference, so I have the above requirements.Jiang
08/09/2023, 10:33 AMthat supports concurrent inferenceHow does it support concurrent inference?
Jiang
08/09/2023, 10:34 AMconcurrentIt should have one of API bellow: 1. async method/function 2. threading 3. multiprocessing 4. batch inference
hao chen
08/09/2023, 10:40 AMmodel_config = {"device_map": "auto", "torch_dtype": torch.float16, "trust_remote_code": True}
model = transformers.AutoModelForCausalLM.from_pretrained(model_path, **model_config)
I can concurrently call the model. generate method in multiple threads, and it works properly. Now that I put model. generate in the runner method, its concurrency advantage is gone.hao chen
08/09/2023, 10:43 AM@bentoml.Runnable.method()
def stream_generate_1(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)
@bentoml.Runnable.method()
def stream_generate_2(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)
@bentoml.Runnable.method()
def stream_generate_3(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)
@bentoml.Runnable.method()
def stream_generate_4(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)Jiang
08/09/2023, 10:44 AMconcurrently call the model. generate method in multiple threads,Something like this?
def generate_output():
input_ids = tokenizer.encode(text, return_tensors='pt')
output = model.generate(input_ids)
generated_text = tokenizer.decode(output[0], skip_special_tokens=True)
print(generated_text)
num_threads = 4
threads = []
for _ in range(num_threads):
thread = threading.Thread(target=generate_output)
thread.start()
threads.append(thread)hao chen
08/09/2023, 10:45 AMJiang
08/09/2023, 10:55 AMAutoModelForCausalLM
We know it is a PyTorch model, right?
Thus one single generate call will spread to multiple threads(same as CPU cores) by PyTorch, internally. https://pytorch.org/docs/stable/notes/cpu_threading_torchscript_inference.html
PyTorch is very smart and already did that for you.
There're two concerns to use threading manually like this:
1. PyTorch is thread-safe. But there is no guarantee for most pre-trained models to be thread-safe
2. Since each generate will already spread N threads, with above script, there will be NxN threads finally on the system. Given that threads:cores = 1:1 give us the best performance, having NxN threads typically takes worse performance.hao chen
08/09/2023, 11:00 AMJiang
08/09/2023, 11:08 AMJiang
08/09/2023, 11:10 AMhao chen
08/09/2023, 11:11 AMreturn await anyio.to_thread.run_sync(
functools.partial(method, **kwargs),
*args,
#limiter=self._limiter,
)
Another issue is the core dispatch, which restricts a runner method from being called concurrently, which is also why I create multiple runner methods.
@bentoml.Runnable.method()
def stream_generate_1(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)
@bentoml.Runnable.method()
def stream_generate_2(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)
@bentoml.Runnable.method()
def stream_generate_3(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)
@bentoml.Runnable.method()
def stream_generate_4(self, inp: List[str], ws_urls: List[str]) -> List[str]:
return self._stream_generate(inp, ws_urls)Jiang
08/09/2023, 11:14 AMhao chen
08/09/2023, 11:20 AMJiang
08/09/2023, 11:26 AMhao chen
08/09/2023, 11:26 AMclass RunnerMethods():
def __init__(self, runner: bentoml.Runner):
self.runner = runner
self.stream_methods, self.non_stream_methods = self.get_all_methods()
self.stream_next_index = -1
self.non_stream_next_index = -1
self.lock = threading.Lock()
print("stream method count: ", len(self.stream_methods))
print("non stream method count: ", len(self.non_stream_methods))
def get_all_methods(self):
stream_methods = []
non_stream_methods = []
for method in self.runner.runner_methods:
if method.name.startswith("stream"):
stream_methods.append(method)
else:
non_stream_methods.append(method)
return stream_methods, non_stream_methods
def get_next_stream_method(self):
with self.lock:
self.stream_next_index += 1
if self.stream_next_index >= len(self.stream_methods):
self.stream_next_index = 0
return self.stream_methods[self.stream_next_index]
def get_next_non_stream_method(self):
with self.lock:
self.non_stream_next_index += 1
if self.non_stream_next_index >= len(self.non_stream_methods):
self.non_stream_next_index = 0
return self.non_stream_methods[self.non_stream_next_index]Jiang
08/09/2023, 11:36 AMAaron Pham
08/09/2023, 12:21 PMhao chen
08/09/2023, 12:23 PMJiang
08/09/2023, 12:48 PMgenerate with multithreading. We did encountered the race condition. It seems that the model is not thread-safe.Jiang
08/09/2023, 12:49 PMhao chen
08/09/2023, 12:51 PMJiang
08/09/2023, 12:53 PMhao chen
08/09/2023, 12:56 PMJiang
08/09/2023, 12:58 PM