スマートフォン・タブレットからインターネットサーバーオペレーション

APPW.jp
 

RAG System / consumer.py

低スペックVPSで動かす
RAGエンジンの設計と実装

RabbitMQ キューイング × llama.cpp × ChromaDB による非同期RAGシステム。メモリ2GB・CPU3コアのConoHa VPS環境での実用を前提にした設計判断を解説します。

💾 RAM 2GB ⚙️ CPU 3コア 🐇 RabbitMQ 🦙 llama.cpp 🔍 ChromaDB

01 — Architecture

システム構成図

consumer.py はRAGシステムの推論エンジンです。フロントエンドからのリクエストを直接受け取るのではなく、RabbitMQ キューを経由して非同期で処理します。フロント側の ws_server.py(WebSocket + Producer)と責務を完全に分離することで、推論の遅延がWebSocket接続に影響しない構造になっています。

ws_server.py WebSocket Producer 🐇 RabbitMQ rag_request リクエストキュー rag_reply 返信キュー rag_cancel キャンセルキュー consumer.py on_message() asyncio ループ prefetch_count=1 poll_cancel_queue() cancelled_sessions ChromaDB ベクトル検索 llama.cpp LLM推論 MongoDB クエリログ リクエスト 返信 キャンセル(非同期)

キューは3種類に分かれています。rag_request(クエリ・スクレイピング等のリクエスト)、rag_reply(ストリーミングトークン・結果の返信)、rag_cancel(処理中断の専用チャンネル)です。キャンセルを専用キューに分離することで、重いLLM推論を開始する前に弾く仕組みが成立します。

02 — VPS Design Decisions

低スペックVPSで動かすための3つの設計判断

メモリ2GB・CPU3コアという制約は、LLM推論(llama.cpp)が最優先でリソースを使う前提に立つと、「余計な処理をいかに排除するか」がアーキテクチャの基本方針になります。

🚫 キャンセルキューでムダな推論を排除する メモリ保護

ユーザーがページを閉じる・再送信するといったセッション切断はWebアプリでは頻繁に起きます。問題は、RabbitMQにキューイング済みのリクエストはセッション切断を知る術がないことです。そのままにすると、誰も受け取らない回答のためにLLM推論がメモリを占有します。

consumer.py — キャンセルキュー購読
cancelled_sessions: set[str] = set()

async def poll_cancel_queue(cancel_q):
    async def on_cancel(message: aio_pika.IncomingMessage):
        async with message.process():
            payload    = json.loads(message.body.decode())
            session_id = payload.get("session_id", "")
            cancelled_sessions.add(session_id)   # ← セットに追加するだけ
            print(f"🚫 キャンセル受信: session={session_id}")

    await cancel_q.consume(on_cancel)

# メインループ側:処理開始前にチェック
if session_id in cancelled_sessions:
    print(f"⏭️  スキップ(キャンセル済み): session={session_id}")
    return   # LLM推論に進まずに捨てる

poll_cancel_queue()asyncio.create_task() でメインループと並走する独立タスクです。キャンセルは cancelled_sessions という単純な set に積まれ、on_message() が次のリクエストをデキューした瞬間に確認します。LLM推論が始まる前に弾けるため、最もコストの高い処理を確実に回避できます。

⚙️ prefetch_count=1 で推論を直列化する 過負荷防止

RabbitMQのデフォルトでは、コンシューマーは次々とメッセージをプリフェッチして並列処理しようとします。LLM推論は1回で数百MB〜1GB近くのメモリを消費するため、2件同時に走るとメモリ不足でプロセスごとクラッシュします。

consumer.py — QoS設定
channel = await connection.channel()
await channel.set_qos(prefetch_count=1)   # ← 1件処理完了後に次を取得

set_qos(prefetch_count=1) の1行で、前のリクエストの ack(完了確認)が返るまで次のメッセージを取得しません。これにより推論が必ず直列化され、メモリ使用量が予測可能になります。スループットは下がりますが、低スペック環境では「落ちないこと」が最優先です。

⚠️

RabbitMQのデフォルト(prefetch_count=0)は「無制限」です。明示的に1を指定しないと低スペックVPSでは予期しないクラッシュが起きます。

asyncio.to_thread() でイベントループを止めない CPU効率

SentenceTransformer・ChromaDB・MongoDB・requestsはいずれも同期的(ブロッキング)なAPIです。そのままasyncioのイベントループで呼ぶと、処理が終わるまで他のコルーチン(キャンセルポーリングなど)が一切動けなくなります。

consumer.py — retrieve() 内の非同期化
async def retrieve(question: str) -> list[dict]:

    # 埋め込み計算(SentenceTransformer: ブロッキング)
    def _encode(text):
        return embedder.encode([text], normalize_embeddings=True).tolist()

    q_emb = await asyncio.to_thread(_encode, question)   # ← スレッドへ退避

    # ChromaDBクエリ(ブロッキング)
    def _query():
        return collection.query(
            query_embeddings=q_emb,
            n_results=EMBED["top_k"],
            include=["documents", "metadatas", "distances"],
        )

    res = await asyncio.to_thread(_query)   # ← スレッドへ退避

asyncio.to_thread() は関数をスレッドプール上で実行し、完了をコルーチンとして待てます。これにより、ChromaDBやMongoDBへのI/O待ち時間中もイベントループが動き続け、キャンセルポーリングやトークン送信が止まりません。3コアのCPUをLLM推論スレッドとI/Oスレッドに効率よく分散できます。

03 — RAG Pipeline

RAGパイプライン コード解説

handle_query() が質問を受け取ってから回答を返すまでの処理を4ステップで解説します。

1
ベクトル検索と距離フィルタリング(retrieve)

質問文を埋め込みベクトルに変換し、ChromaDBで類似チャンクを検索します。検索結果は上位 top_k 件ですが、距離閾値でさらに絞り込みます。

consumer.py — 距離フィルタ
contexts = []
for doc, meta, dist in zip(res["documents"][0],
                            res["metadatas"][0],
                            res["distances"][0]):
    if dist < EMBED["distance_threshold"]:   # ← 閾値未満のみ採用
        contexts.append({
            "text":     doc,
            "url":      meta.get("url", ""),
            "title":   meta.get("title", ""),
            "distance": round(dist, 4),
        })

距離が閾値を超えるチャンクは関連度が低いと判断して除外します。contexts が空になった場合は「未カバー」として MongoDB にフラグを立てて早期リターンします。これにより、無関係なコンテキストがLLMに混入して幻覚回答を誘発するのを防いでいます。

2
トークン数管理(trim_contexts)

コンテキストが多い場合、プロンプト全体が n_ctx(コンテキスト長上限)を超えてllama.cppがエラーになります。trim_contexts() はトークナイザーで実際のトークン数を計算しながらコンテキストを切り詰めます。

consumer.py — trim_contexts()
MAX_CONTEXT_TOKENS = LLM["n_ctx"] - LLM["max_tokens"] - 200  # 余裕を持たせる

def trim_contexts(contexts, question):
    base_prompt  = f"質問: {question}\n回答: "
    base_tokens  = len(llm.tokenize(base_prompt.encode()))

    trimmed, used_tokens = [], base_tokens
    for ctx in contexts:
        ctx_tokens = len(llm.tokenize(ctx["text"].encode()))
        if used_tokens + ctx_tokens > MAX_CONTEXT_TOKENS:
            break               # 超えたらそこで打ち止め
        trimmed.append(ctx)
        used_tokens += ctx_tokens
    return trimmed
💡

上限 = n_ctx − max_tokens − 200 という計算でプロンプト側とレスポンス側のトークン予算を明示的に分離しています。余白の200はシステムプロンプトや改行文字のズレを吸収するバッファです。

3
LLMストリーミング:同期APIをasyncioに橋渡し

llama.cppのストリームAPIは同期ジェネレータです。これをasyncioのイベントループから直接呼ぶとスレッドをブロックします。threading.Thread + queue.Queue でブリッジします。

consumer.py — ストリーミングブリッジ
token_queue = queue.Queue()

def _llm_stream(prompt, token_q):
    for llm_chunk in llm(prompt, max_tokens=LLM["max_tokens"],
                       temperature=LLM["temperature"], stream=True):
        token = llm_chunk["choices"][0]["text"]
        token_q.put(token)

        # max_tokens超過の検出
        finish = llm_chunk["choices"][0].get("finish_reason")
        if finish == "length":
            token_q.put("\n\n(※ 回答が長くなったため省略されました)")
    token_q.put(None)   # 終端シグナル

# 別スレッドで推論を開始
thread = threading.Thread(target=_llm_stream, args=(prompt, token_queue), daemon=True)
thread.start()

# asyncio側: run_in_executor で queue.get をノンブロッキングに
loop = asyncio.get_running_loop()
while True:
    token = await loop.run_in_executor(None, token_queue.get)
    if token is None:
        break

推論スレッドはトークンを queue.Queueput() し、asyncio側は run_in_executorqueue.get() をノンブロッキングに待ちます。終端は None をセンチネルとして使います。

4
トークンのバッファリング送信

生成されるトークンを1つずつ即時にRabbitMQへ送ると、メッセージ数が数百件に膨らみオーバーヘッドが大きくなります。そのため、8文字 or 句読点を契機にまとめて送ります。

consumer.py — バッファリング
token_buf  = ""
answer_buf = []   # 全文収集(ログ保存用)

while True:
    token = await loop.run_in_executor(None, token_queue.get)
    if token is None:
        break
    token_buf  += token
    answer_buf.append(token)

    # 8文字 OR 句読点でフラッシュ
    if len(token_buf) >= 8 or token in ("。", "、", "\n"):
        await send({"type": "token", "content": token_buf})
        token_buf = ""

# 残りを送出
if token_buf:
    await send({"type": "token", "content": token_buf})

answer_buf は別の目的(ログ保存)のために全トークンを収集しています。バッファを flush するたびに answer_buf には追記され、推論完了後に "".join(answer_buf) で回答全文が得られます。

04 — Scraping & Chunk Strategy

セマンティックDOMスクレイピングとチャンク戦略

RAGの回答品質はチャンクの質に直結します。1ページをひとつの巨大なテキストとして扱うよりも、意味のある単位(用語・セクション)でチャンクを分けたほうがベクトル検索の精度が上がり、レスポンスも速くなることが実験でわかりました。そのため、ページの構造によってスクレイピング戦略を2段階に分けています。

優先順位付きDOM選択

新規・修正ページ
article.term-block article.chunk-block 以降は既存ページと同順
既存ページ
div.asset-content article div.contentsWrapper main body
A
term-block / chunk-block:意味境界での事前分割

新規・修正ページには <article class="term-block"> または <article class="chunk-block"> 単位で用語解説やセクションをマークアップしています。これが見つかった場合、各ブロックを \n\n---\n\n で結合してテキストを作ります。

consumer.py — term-block 優先取得
term_blocks = soup.find_all("article", class_="term-block")

if not term_blocks:
    term_blocks = soup.find_all("article", class_="chunk-block")

if term_blocks:
    text = "\n\n---\n\n".join(
        b.get_text(separator="\n", strip=True)
        for b in term_blocks
    )   # セパレータがチャンク分割の境界になる

\n\n---\n\n は後述の chunk_text() がセパレータとして認識します。これにより、「1用語 = 1チャンク」に近い粒度でChromaDBに登録されます。

B
二段階チャンク分割(セパレータ優先 → 語数フォールバック)
consumer.py — chunk_text()
def chunk_text(text: str) -> list[str]:
    if "\n\n---\n\n" in text:
        blocks = text.split("\n\n---\n\n")
        chunks = []
        for block in blocks:
            if len(block.split()) <= EMBED["chunk_size"]:
                if block.strip():
                    chunks.append(block.strip())   # 1ブロック = 1チャンク
            else:
                chunks.extend(_word_chunk(block))  # 長すぎれば再分割
        return chunks

    return _word_chunk(text)   # セパレータなし → 語数ベース


def _word_chunk(text: str) -> list[str]:
    words   = text.split()
    size    = EMBED["chunk_size"]
    overlap = EMBED["chunk_overlap"]
    chunks, i = [], 0
    while i < len(words):
        chunk = " ".join(words[i:i + size])
        if chunk.strip():
            chunks.append(chunk)
        i += size - overlap   # オーバーラップ分だけ戻る
    return chunks

_word_chunk()overlap は前のチャンクの末尾N語を次のチャンク先頭に重複させます。チャンク境界で文脈が途切れるのを防ぐための工夫です。

C
差分チェックによる埋め込み再計算のスキップ

スクレイピング済みのURLが再登録されるとき、本文が変わっていなければ埋め込み計算とupsertをスキップします。センテンストランスフォーマーの推論はCPUで数秒かかるため、不変コンテンツへのコスト削減として重要です。

consumer.py — 差分チェック
# URLに紐づく既存チャンクを一括取得(1URLにつき1回のAPI呼び出し)
existing     = collection.get(where={"url": url})
existing_ids = existing["ids"]
existing_docs = {
    cid: doc for cid, doc in zip(existing_ids, existing["documents"])
}

# IDリストとテキスト内容が両方完全一致なら変更なし
is_unchanged = (
    existing_ids == new_cids
    and all(existing_docs.get(cid) == text
           for cid, text in zip(new_cids, chunks))
)
if is_unchanged:
    continue   # 埋め込み計算・upsert をスキップ

「IDが一致 かつ テキストが全チャンク一致」の条件です。1チャンクでも変わっていれば古いチャンクを全削除してから作り直します。部分更新より全削除・全再登録のほうが一貫性を保ちやすいためです。

05 — Logging & Observability

ロギングと観測性

「動かして終わり」ではなく、「どの質問が答えられないか」「どのカテゴリが多いか」「どこが遅いか」をデータで把握するための仕組みです。

バックグラウンドログ保存(回答速度に影響させない)

カテゴリ分類はLLM推論を1回追加で実行するため数秒かかります。これをユーザーへの返信前に行うとレスポンスが遅くなります。asyncio.create_task() で返信後にバックグラウンド実行することでレスポンスタイムへの影響をゼロにしています。

consumer.py — バックグラウンドログ保存
# 回答送信を完了させてから非同期でログを保存
await send({"type": "answer", "content": "", "sources": sources})

async def _log_after_reply():
    category = await asyncio.to_thread(classify_question, question)
    await asyncio.to_thread(
        save_query_log,
        session_id   = session_id,
        question     = question,
        answer       = answer_text,
        contexts     = contexts,
        elapsed_sec  = elapsed,
        is_uncovered = is_uncovered,
        category     = category,
    )

asyncio.create_task(_log_after_reply())  # ← 返信後にバックグラウンド実行

LLMによる質問カテゴリ自動分類

consumer.py — classify_question()
CATEGORIES = ["FP用語", "財務・会計", "IT・技術", "ブロックチェーン", "その他"]

def classify_question(question: str) -> str:
    result = llm(
        f"以下の質問を1つのカテゴリに分類してください。\n"
        f"カテゴリ: {' / '.join(CATEGORIES)}\n"
        f"質問: {question}\n"
        f"カテゴリのみ回答(他の文字は不要):",
        max_tokens=20,
        temperature=0.0,   # ← 決定論的(毎回同じ結果)
        stop=["\n", "。", "、"],
    )
    raw = result["choices"][0]["text"].strip()
    for cat in CATEGORIES:
        if cat in raw:
            return cat
    return "その他"
💡

temperature=0.0 は完全に決定論的な出力にします。分類タスクは創造性が不要なため、毎回同じ入力に同じカテゴリが返る安定動作が重要です。max_tokens=20 で応答を短く打ち切り、推論コストを最小化しています。

MongoDBクエリログのスキーマと4モード取得

consumer.py — ログドキュメント構造
doc = {
    "timestamp":    datetime.now(timezone.utc),
    "session_id":   session_id,
    "question":     question,
    "answer":       answer,
    "sources":      list({c["url"] for c in contexts}),  # 重複排除
    "chunks_used":  len(contexts),
    "category":     category,          # LLM分類結果
    "is_uncovered": is_uncovered,      # 未カバー質問フラグ
    "is_error":     is_error,
    "elapsed_sec":  round(elapsed, 2), # 処理時間(秒)
}

handle_get_log()tab パラメータで取得モードを切り替えます。

tab値取得内容ソート
recent直近のクエリ(デフォルト)timestamp 降順
uncoveredis_uncovered: true のクエリ
→ 知識ベースに登録すべき未カバー質問の洗い出し
timestamp 降順
slow処理時間が長いクエリ
→ パフォーマンスのボトルネック特定
elapsed_sec 降順
erroris_error: true のクエリtimestamp 降順

MongoDBフォールトトレランス

consumer.py — 接続失敗時の継続動作
try:
    mongo_client = MongoClient("mongodb://localhost:27017/")
    log_col      = mongo_client["rag_system"]["query_log"]
    print("✅ [Consumer] MongoDB 接続完了")
except Exception as e:
    print(f"⚠️  MongoDB 接続失敗(ログなしで継続): {e}")
    log_col = None   # ← None にしておき save_query_log で早期リターン

def save_query_log(...) -> None:
    if log_col is None:   # ← 接続失敗時は何もしない
        return

MongoDBが落ちていてもRAG本体(検索・推論・返信)は継続動作します。ログDBの障害がユーザー向けサービスに連鎖しない設計です。

06 — Action Reference

アクション一覧リファレンス

on_message() が受け取るペイロードの action フィールドで処理を振り分けます。全8アクションの概要と返送メッセージタイプをまとめます。

action処理内容主な返送 type
query 質問受信 → ベクトル検索 → LLM推論(ストリーミング) status / context_preview / token / answer
scrape URL配列を受け取りスクレイピング → ChromaDB登録 status / answer
db_info ChromaDBの登録チャンク数・記事一覧を返す answer
list_articles 登録済み記事をURL単位で集約してチャンク数付きで返す article_list
delete_article 指定URLの全チャンクをChromaDBから削除 delete_result / error
rescrape 指定URL:既存チャンク削除 → 再スクレイピング → 再登録 status / rescrape_result
rescrape_all 全登録URL順次再スクレイピング(間隔あり) status / rescrape_all_result
get_log MongoDBからクエリログ取得(tab別モード) query_log / error
💡

全アクションのペイロードには session_id(クライアント識別)と correlation_id(リクエスト識別)が必須です。consumer.py はすべての返送メッセージにこの2つを付与して返すため、複数クライアントが同時接続していても返信の宛先を正しく判別できます。

『consumer.py — RAGエンジン解説』を公開しました。