This message was deleted.
# ask-for-help
s
This message was deleted.
🍱 2
👍 1
j
Hi chen. Great finding. Runner is designed for computing bound workloads, mostly model inference on CPU or GPU. All of them have internal multi-threading mechanic and can utilize all the cores, and also a blocking sync API. In this case async concurrency doesn't makes sense, since it will be blocked by the blocking API of ML frameworks.
Because of that, we introduced batching optimization for runners. It is an approach that can also implement concurrency, and fit better for models.
If your workload is IO bound and asyncio can help, leave them to the API server side (svc.api), and define api methods in async flavor.
Copy code
@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]
h
Thank you for your reply. I now have a batch runner method:
Copy code
@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.
j
I see.
that supports concurrent inference
How does it support concurrent inference?
concurrent
It should have one of API bellow: 1. async method/function 2. threading 3. multiprocessing 4. batch inference
h
The following model instance is created in this way, and its generate method supports concurrency.
Copy code
model_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.
In order to achieve better concurrency performance, I have to write the code as follows😂:
Copy code
@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)
j
concurrently call the model. generate method in multiple threads,
Something like this?
Copy code
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)
h
right👍
j
Copy code
AutoModelForCausalLM
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.
h
Yes, when model. generate is concurrent with too many threads, the speed of token generation under a single thread will decrease very quickly, but this can be achieved by limiting the number of concurrent threads to achieve an acceptable speed of token generation under a single thread. The runner method would be perfect if it supports configurable concurrency.
j
We can but it is not what ML frameworks suggested to be used. And from our test its benefit doesn't worth the risk of possible thread race. But we are open for that, model inferences are different from each other. Would you mind sharing us your test result, if you find a greater value give us better performance? @hao chen
h
I found that there are two limitations to concurrency in bentoml in the runner. One is what I mentioned earlier, which limits the runner to only call one runner method of that runner at the same time (if the runner defines multiple runner methods). Now I have commented out the code there.
Copy code
return 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.
Copy code
@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)
j
I mean it is designed like this. Did you get better overall throughput for a long running with your tweak?
h
I only have one GPU and there may be multiple people accessing my inference service. If there is no concurrency, one person must wait for the previous person's inference request to end before starting their inference request, which is not very good in terms of experience. If concurrency is supported, multiple people can receive the results returned by the inference service at the same time. Although the rate at which each person receives the results through stream has slowed down, it is acceptable as long as it is not too slow.
j
I see. It is a valid one🙂.
h
Now I choose which runner method to use in sequence in the svc. api. From the perspective of the effect, it is acceptable, and the response speed of three concurrent requests to the inference service is also acceptable.
Copy code
class 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]
j
Yes. I don't see good workaround here. Will investigate the LLM case, it seems that we should support it from the design. Thanks for your contribution!
a
Hey, what is the model that you are using?
h
the model is Baichuan 13b,a pre-train LLM model
j
@hao chen We tried bauchuan 13b
generate
with multithreading. We did encountered the race condition. It seems that the model is not thread-safe.
Using multi-threading over it may cause unexpected errors
h
Oh, this is something I didn't expect. I just did a simple concurrency test before. May I ask how many threads you used
j
Only four, and they have different prompts
h
Thank you very much. I think I need to do more concurrent testing
j
🍻 Also appreciate for your great question.