🙆‍♀️

vLLMつかったら、ベクトル検索サーバーが高速化して安定性も向上した

に公開

この記事は Uzabase Advent Calendar 2025 の 18 日目の記事です。

https://qiita.com/advent-calendar/2025/uzabase

はじめに

プロダクトチームでは、Speeda AI Agentを開発していますが、端々でベクトル検索を用いています。また、ベクトル検索するためのモデル開発や、そのホスティングも自前でおこなっています。

ベクトル検索用のモデルは社内ファインチューニングしているため、社内のベンチマークデータでは、Geminiのベクトルを精度で大きく上回っており、ユーザー体験に貢献していると考えています。また、基盤も安定しており、ニュースのような頻繁かつ断続的に追加されるデータも安定してさばけるような基盤を構築しています。

本シリーズではここに至るまでの以下の技術的な工夫を紹介して行きたいと思います。

  • vLLM化による速度と安定性の向上
  • HPAでスケールできるインフラの構築
  • Sentence Transformersで省計算資源で学習を行う工夫

今回は、「vLLM化による速度と安定性の向上」について記載します。

直面した課題

初期のころは、fastapiを用いて、自分たちで作成したサーバーでデプロイをしていました。しかし、徐々に以下の問題が明らかになってきました。

  1. 速度不足:定期的に大量のデータをベクトル化する必要があるニュースなどのデータにおいて、1サーバーあたりの処理速度が足らず、複数のサーバーを使うことになりました。プロダクトチームではGCPを利用しており、サーバーの費用がかなり高額になってしまいました。
  2. 安定性不足:多くのリクエストを受け付けると、リクエストを失敗するということが頻発しており、クライアント側でリトライを何回かかけてもらうなどの処理を行っていました。そのため、1と合わせて、リクエストを処理しきれない場面があるのではないかという不安がありました。

これらの対策として、vLLMを用いることにしました。

vLLM

vLLMは、LLMの推論・サービングに特化したフレームワークです。特に、大量のリクエストを捌くという点において優れています。今回は、vLLMの特徴のうち、速度不足と安定性不足に寄与した特徴を紹介し、特にDynamic Batchingについて実装を見ます。

速度不足解消に寄与した特徴

速度不足に効いた工夫は、flash-attentionです。flash-attentionはattention計算を高速化して、メモリを削減する計算手法です。Attentionの行列計算を小分けにして行うことで、O(n^2)の計算量を回避することでメモリ削減をしています。また、高速なSRAM領域をつかったりAttentionの演算をまとめたカーネルを用いることで高速化しています。概念図は以下の通りです。


公式githubより引用

詳細な解説はネットにいっぱいありそうなので割愛します。また、同時にGPUをT4からL4に変更しており、手元集計では合わせて大体10倍ぐらいのスループットが出ていました。flash-attention自体は、元々のお手製サーバーでも可能でしたが、vLLMではコンテナとして配布されているため、自らbuildする必要がないという点で、パイプラインがシンプルになるという嬉しさもありました。

安定性不足解消に効いた特徴

Dynamic Batchingが挙げられます。これは複数のリクエストをまとめるものです。これによって、テキストをベクトル化する際の計算負荷をクライアントから切り離して、比較的均一なGPU負荷にすることが可能です。詳細は以下の記事を参照ください。

https://zenn.dev/rinna/articles/7d10e61f694611#dynamic-batching-(request-level-scheduling)

なお、vLLMにはさらに高度なContinuous Batchingが実装されていますが、ベクトル検索に使用する場合encode(リンク先記事ではprefill)しかないため、Dynamic Batchingとほとんど同等になると思われます(本当はChunked Prefill次第ですが調査しきれず)。Continuous Batchingの詳細はこちらが参考になります。
https://huggingface.co/blog/continuous_batching

それでは、実装を見ていきます。

vLLMのアーキテクチャ

Continuous Batchingは複数のクラスが連動する形で実装されているため、まずは全体像としてアーキテクチャを見ます。vLLMのアーキテクチャは以下のようになっています。下図のLLM Engine内にContinuous Batchingは実装されています。


公式ドキュメントより引用

このアーキテクチャ図を見ると、APIもpythonのLLM classもどちらもLLM Engineをつかっていることがわかります。そのため、わかりやすさを優先し、LLM classの実装からみていきます。

vLLMはどのようにモデルを実行しているのか

LLMクラスでテキストをエンコードするときは、LLM.embed()を使用します。ざっとモデルの出力を得られるまでの流れは以下です。

1. LLM.embed()
2. LLM.encode()
3. LLMEngine.step()
4. InprocClient.get_output()
5. EngineCore.step_fn()
6. EngineCore.step() または step_with_batch_queue()

EngineCore.step_with_batch_queue()はパイプライン並列化時に使用するようです。そのため、今回はEngineCore.step()の中身をみます。EngineCore.step()v1/engine/core.py内にあり以下の通りです。

class EngineCore:
    def __init__(...):
        ...
        self.model_executor = executor_class(vllm_config)
        ...
        self.scheduler: SchedulerInterface = Scheduler(...)
        ...
        self.step_fn = (
            self.step if self.batch_queue is None else self.step_with_batch_queue
        )
        ...
        
    def step(self) -> tuple[dict[int, EngineCoreOutputs], bool]:
        ,..
        scheduler_output = self.scheduler.schedule()
        future = self.model_executor.execute_model(scheduler_output, non_block=True)
        grammar_output = self.scheduler.get_grammar_bitmask(scheduler_output)
        with self.log_error_detail(scheduler_output):
            model_output = future.result()
            if model_output is None:
                model_output = self.model_executor.sample_tokens(grammar_output)

        engine_core_outputs = self.scheduler.update_from_output(
            scheduler_output, model_output
        )

        return engine_core_outputs, scheduler_output.total_num_scheduled_tokens > 0
    ...

scheduler.schedule()からの出力がmodel_executor.execute_model()に渡されて、モデルが実行されています。その出力であるmodel_outputを再度scheduler.update_from_output()に渡してschedulerを更新しています。ということで、schdulerがリクエストを作っていることが推察できます。

Schedulerについて

先程までで、scheduler.schedule()がdynamic batchingのためのリクエスト作成をしており、scheduler.update_from_output()がモデルの処理結果から状態の更新を行っていると推察されます。そのため、これの2つを見ていきます。本節でみる実装は、vllm/v1/core/sched/scheduler.pyにあります。

まずSchedulerクラスについて詳細を見ていきます。Schedulerはwaitingキューとrunningキューを持っており、waitingキューに新規リクエストを取得が追加され、runningキューに継続リクエストが入っています。

class Scheduler():
    def __init__(...) -> None:
        ...
        self.waiting = create_request_queue(self.policy)
        self.running: list[Request] = []
        ...

LLM.embed()をエントリーポイントとする場合に、waitingキューにリクエストが入るまでのながれは以下の通りです。

1. LLM.embed() → LLM.encode()を呼び出し
2. LLM.encode() → _validate_and_add_requests()で各プロンプトを処理
3. LLM._add_request() → _process_inputs()でプロンプトをトークン化しEngineCoreRequestを作成
4. LLMEngine.add_request() → OutputProcessor.add_request()で出力処理用に登録
5. LLMEngine.add_request() → EngineCoreClient.add_request()を呼び出し
6. EngineCoreClient.add_request() -> EngineCoreRequestクラスをRequestクラスに変換し、EngineCore.add_request()を呼び出し
- LLMクラスからの場合は、実体はInprocClientクラス
7. EngineCore.add_request() → バリデーション後、Scheduler.add_request()を呼び出し
8. Scheduler.add_request() → waiting.add_request()で待機キューに追加し、self.requestsに登録

Schedulerへは、Requestクラスのオブジェクトが登録されます。入力からの変換は、まずLLM._add_request()内で、self._preprocess_inputs()によって、EngineCoreRequestクラスに変換され、LLMEngine.add_request()に渡されます。その後、EngineCoreClient.add_request() 内でRequestクラスに変換されます。Requestクラスには、request_idやprompt_token_idsなどが含まれています。名前から推察できますが、prompt_token_idsはすでにtokenizeがなされてtoken_idの列になっていると予想されます。tokenizeは、self._preprocess_inputs()内で行われています。この処理は、LLMEngine.input_processor()にて行われています。

Scheduler.schedule()について

さて、肝心のScheduler.schedule()の中身ですが以下のようになっています。

class Scheduler(ShedulerInterface):
    def schedule(self) -> SchedulerOutput:
		# 1. prepare containers of requests
        scheduled_new_reqs: list[Request] = []
        scheduled_resumed_reqs: list[Request] = []
        scheduled_running_reqs: list[Request] = []
        preempted_reqs: list[Request] = []
        
        req_to_new_blocks: dict[str, KVCacheBlocks] = {}
        num_scheduled_tokens: dict[str, int] = {}
        token_budget = self.max_num_scheduled_tokens
        encoder_compute_budget = self.max_num_encoder_input_tokens
        ...
        # 2. schedule the RUNNING requests.
        req_index = 0
        while req_index < len(self.running) and token_budget > 0:
            request = self.running[req_index]
            ...
            num_new_tokens = (
                request.num_tokens_with_spec
                + request.num_output_placeholders
                - request.num_computed_tokens
            )
            ...
            scheduled_running_reqs.append(request)
            req_to_new_blocks[request.request_id] = new_blocks
            num_scheduled_tokens[request.request_id] = num_new_tokens
            token_budget -= num_new_tokens
            req_index += 1
        ...
        # 3. schedule the WAITING requests.
        if not preempted_reqs:
            while self.waiting and token_budget > 0:
                if len(self.running) == self.max_num_running_reqs:
                    break

                request = self.waiting.peek_request()
                ...
                if load_kv_async:
                    # KVTransfer: loading remote KV, do not allocate for new work.
                    assert num_external_computed_tokens > 0
                    num_new_tokens = 0
                else:
                    # Number of tokens to be scheduled.
                    # We use `request.num_tokens` instead of
                    # `request.num_prompt_tokens` to consider the resumed
                    # requests, which have output tokens.
                    num_new_tokens = request.num_tokens - num_computed_tokens
                ...
                req_index += 1
                self.running.append(request)
                if self.log_stats:
                    request.record_event(
                        EngineCoreEventType.SCHEDULED, scheduled_timestamp
                    )
                if request.status == RequestStatus.WAITING:
                    scheduled_new_reqs.append(request)
                elif request.status == RequestStatus.PREEMPTED:
                    scheduled_resumed_reqs.append(request)
                else:
                    raise RuntimeError(f"Invalid request status: {request.status}")
                ...
                num_scheduled_tokens[request.request_id] = num_new_tokens
                token_budget -= num_new_tokens
                request.status = RequestStatus.RUNNING
                request.num_computed_tokens = num_computed_tokens
                
        # 4. Construct the scheduler output.
        if self.use_v2_model_runner:
            scheduled_new_reqs = scheduled_new_reqs + scheduled_resumed_reqs
            scheduled_resumed_reqs = []
            new_reqs_data = [
                NewRequestData.from_request(
                    req,
                    req_to_new_blocks[req.request_id].get_block_ids(),
                    req._all_token_ids,
                )
                for req in scheduled_new_reqs
            ]
        else:
            new_reqs_data = [
                NewRequestData.from_request(
                    req, req_to_new_blocks[req.request_id].get_block_ids()
                )
                for req in scheduled_new_reqs
            ]

        with record_function_or_nullcontext("schedule: make_cached_request_data"):
            cached_reqs_data = self._make_cached_request_data(
                scheduled_running_reqs,
                scheduled_resumed_reqs,
                num_scheduled_tokens,
                scheduled_spec_decode_tokens,
                req_to_new_blocks,
            )

        # Record the request ids that were scheduled in this step.
        self.prev_step_scheduled_req_ids.clear()
        self.prev_step_scheduled_req_ids.update(num_scheduled_tokens.keys())
        
        # 5. materialize output       
        scheduler_output = SchedulerOutput(
            scheduled_new_reqs=new_reqs_data,
            scheduled_cached_reqs=cached_reqs_data,
            num_scheduled_tokens=num_scheduled_tokens,
            total_num_scheduled_tokens=total_num_scheduled_tokens,
            scheduled_spec_decode_tokens=scheduled_spec_decode_tokens,
            scheduled_encoder_inputs=scheduled_encoder_inputs,
            num_common_prefix_blocks=num_common_prefix_blocks,
            preempted_req_ids={req.request_id for req in preempted_reqs},
            # finished_req_ids is an existing state in the scheduler,
            # instead of being newly scheduled in this step.
            # It contains the request IDs that are finished in between
            # the previous and the current steps.
            finished_req_ids=self.finished_req_ids,
            free_encoder_mm_hashes=self.encoder_cache_manager.get_freed_mm_hashes(),
        )
        with record_function_or_nullcontext("schedule: update_after_schedule"):
            self._update_after_schedule(scheduler_output)
        return scheduler_output

非常に長いコードなので、重要箇所を抜粋しています。この処理でやっていることを大まかに言えば、token_budgetに達するまでrequestを詰め込んで、一つのoutputとして返すということです。それでは、詳細を見ていきます。大まかなブロック毎にコメントをしています。処理の流れとしては以下の通りです。

  • #1.でrequestを詰め込むためのlistやdictを用意しています。
  • #2.で現在実行中(running) のrequestを取り出し、scheduled_running_reqsにappendしています。そしてrequestのtoken数を計算し、token_budgetをそのrequestの分だけ減らしています。
  • #3.では実行していない(waiting)requestを取り出し、scheduled_new_reqsにappendしています。PREEMPTEDの場合はscheduled_resumed_reqsにappendしています。そして、requestのtoken数分をtoken_budgetから減らしています。#2, #3でtoken_budgetを消費するか、それぞれrunningキューとwaitingキューが終端になるかしたら、ループを抜け出します。
  • #4.でoutputに変換するための前処理をします。scheduled_new_reqsはほとんど詰め替えだけですが、scheduled_running_reqsとscheduled_resumed_reqsは_make_cached_request_dataは一つにまとめられます。
  • 最後に#5.でoutputとしてまとめています。

EngineCore.step()を思い出すと、ここでまとめられたoutputがmodel_executor.execute_model()にわたされています。よって、以上の処理で複数のrequestをまとめてモデルに入力するDynamic batchingになっていると考えられます。

Scheduler.update_from_output()について

こちらは、Dynamnic Batchingとは直接的な関係がない &力尽きてきた ため、簡単に述べます。この関数の主な役割は、状態の更新とoutputの詰替えです

    def update_from_output(
        self,
        scheduler_output: SchedulerOutput,
        model_runner_output: ModelRunnerOutput,
    ) -> dict[int, EngineCoreOutputs]:
        sampled_token_ids = model_runner_output.sampled_token_ids
        logprobs = model_runner_output.logprobs
        prompt_logprobs_dict = model_runner_output.prompt_logprobs_dict
        num_scheduled_tokens = scheduler_output.num_scheduled_tokens
        pooler_outputs = model_runner_output.pooler_output
        num_nans_in_logits = model_runner_output.num_nans_in_logits
        kv_connector_output = model_runner_output.kv_connector_output
        ...
        for req_id, num_tokens_scheduled in num_scheduled_tokens.items():
            ...
            prompt_logprobs_tensors = prompt_logprobs_dict.get(req_id)
            if new_token_ids or pooler_output is not None or kv_transfer_params:
                # Add EngineCoreOutput for this Request.
                outputs[request.client_index].append(
                    EngineCoreOutput(
                        request_id=req_id,
                        new_token_ids=new_token_ids,
                        finish_reason=request.get_finished_reason(),
                        new_logprobs=new_logprobs,
                        new_prompt_logprobs_tensors=prompt_logprobs_tensors,
                        pooling_output=pooler_output,
                        stop_reason=request.stop_reason,
                        events=request.take_events(),
                        kv_transfer_params=kv_transfer_params,
                        trace_headers=request.trace_headers,
                        num_cached_tokens=request.num_cached_tokens,
                        num_nans_in_logits=request.num_nans_in_logits,
                    )
                )

このように、model_runner_outputの中身を取り出して、request_id毎に、EngineCoreOutputに詰め替えをしています。
また続きをみると、以下のようにstop状態になったrequestをqueueから削除し、状態を更新しています。

        # Remove the stopped requests from the running and waiting queues.
        if stopped_running_reqs:
            self.running = remove_all(self.running, stopped_running_reqs)
        if stopped_preempted_reqs:
            # This is a rare case and unlikely to impact performance.
            self.waiting.remove_requests(stopped_preempted_reqs)

終わりに

このへんで力尽きたので、今回はここで終わりにしようと思います。Dynamic Batchingの肝の部分については、記載することができたのではないかと思います。なお、もう少し全体が気になるという方は、以下の記事が参考になると思います。

https://nttdocomo-developers.jp/entry/2024/12/19/090000_6

今後機会があれば、Chunked prefillやzmqをどのように使っているかといった部分や、Model Executionの内部やCUDA Graphの部分についても調べてみたいと思います。

Discussion