01 — Architecture
システム構成図
consumer.py はRAGシステムの推論エンジンです。フロントエンドからのリクエストを直接受け取るのではなく、RabbitMQ キューを経由して非同期で処理します。フロント側の ws_server.py(WebSocket + Producer)と責務を完全に分離することで、推論の遅延がWebSocket接続に影響しない構造になっています。
キューは3種類に分かれています。rag_request(クエリ・スクレイピング等のリクエスト)、rag_reply(ストリーミングトークン・結果の返信)、rag_cancel(処理中断の専用チャンネル)です。キャンセルを専用キューに分離することで、重いLLM推論を開始する前に弾く仕組みが成立します。
02 — VPS Design Decisions
低スペックVPSで動かすための3つの設計判断
メモリ2GB・CPU3コアという制約は、LLM推論(llama.cpp)が最優先でリソースを使う前提に立つと、「余計な処理をいかに排除するか」がアーキテクチャの基本方針になります。
ユーザーがページを閉じる・再送信するといったセッション切断はWebアプリでは頻繁に起きます。問題は、RabbitMQにキューイング済みのリクエストはセッション切断を知る術がないことです。そのままにすると、誰も受け取らない回答のためにLLM推論がメモリを占有します。
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推論が始まる前に弾けるため、最もコストの高い処理を確実に回避できます。
RabbitMQのデフォルトでは、コンシューマーは次々とメッセージをプリフェッチして並列処理しようとします。LLM推論は1回で数百MB〜1GB近くのメモリを消費するため、2件同時に走るとメモリ不足でプロセスごとクラッシュします。
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では予期しないクラッシュが起きます。
SentenceTransformer・ChromaDB・MongoDB・requestsはいずれも同期的(ブロッキング)なAPIです。そのままasyncioのイベントループで呼ぶと、処理が終わるまで他のコルーチン(キャンセルポーリングなど)が一切動けなくなります。
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ステップで解説します。
質問文を埋め込みベクトルに変換し、ChromaDBで類似チャンクを検索します。検索結果は上位 top_k 件ですが、距離閾値でさらに絞り込みます。
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に混入して幻覚回答を誘発するのを防いでいます。
コンテキストが多い場合、プロンプト全体が n_ctx(コンテキスト長上限)を超えてllama.cppがエラーになります。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はシステムプロンプトや改行文字のズレを吸収するバッファです。
llama.cppのストリームAPIは同期ジェネレータです。これをasyncioのイベントループから直接呼ぶとスレッドをブロックします。threading.Thread + queue.Queue でブリッジします。
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.Queue に put() し、asyncio側は run_in_executor で queue.get() をノンブロッキングに待ちます。終端は None をセンチネルとして使います。
生成されるトークンを1つずつ即時にRabbitMQへ送ると、メッセージ数が数百件に膨らみオーバーヘッドが大きくなります。そのため、8文字 or 句読点を契機にまとめて送ります。
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 class="term-block"> または <article class="chunk-block"> 単位で用語解説やセクションをマークアップしています。これが見つかった場合、各ブロックを \n\n---\n\n で結合してテキストを作ります。
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に登録されます。
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語を次のチャンク先頭に重複させます。チャンク境界で文脈が途切れるのを防ぐための工夫です。
スクレイピング済みのURLが再登録されるとき、本文が変わっていなければ埋め込み計算とupsertをスキップします。センテンストランスフォーマーの推論はCPUで数秒かかるため、不変コンテンツへのコスト削減として重要です。
# 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() で返信後にバックグラウンド実行することでレスポンスタイムへの影響をゼロにしています。
# 回答送信を完了させてから非同期でログを保存 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による質問カテゴリ自動分類
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モード取得
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 降順 |
| uncovered | is_uncovered: true のクエリ → 知識ベースに登録すべき未カバー質問の洗い出し | timestamp 降順 |
| slow | 処理時間が長いクエリ → パフォーマンスのボトルネック特定 | elapsed_sec 降順 |
| error | is_error: true のクエリ | timestamp 降順 |
MongoDBフォールトトレランス
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つを付与して返すため、複数クライアントが同時接続していても返信の宛先を正しく判別できます。