🔄

Stateful MCPサーバーで社内データ分析エージェントを構築する

に公開

Stateful MCPサーバーで社内データ分析エージェントを構築する

この記事でわかること

  • MCP(Model Context Protocol)のセッション管理とTasks拡張を活用したツール実行状態の永続化設計
  • SQLiteによるチェックポイント永続化と、再接続時の状態復元を実現する具体的なPython実装
  • 社内データ分析エージェントにおける長時間実行クエリの非同期管理パターン
  • Streamable HTTP環境でのMcp-Session-Id管理と再接続ハンドリングの実装
  • 本番環境で運用するための3層状態管理アーキテクチャ(インメモリ / SQLite / Redis)

対象読者

  • 想定読者: 中級〜上級のPythonバックエンドエンジニア
  • 必要な前提知識:
    • Python 3.12+ の非同期プログラミング(asyncio, async/await
    • MCP(Model Context Protocol)の基本概念(ツール、リソース、プロンプト)
    • SQLiteまたはPostgreSQLの基本操作
    • LLMエージェントフレームワーク(LangChain / LangGraph 等)の利用経験

結論・成果

Stateful MCPサーバーを導入することで、社内データ分析エージェントにおいて以下の成果が報告されています。

  • タスク完了率の改善: 長時間実行エージェント(4時間超)では、状態永続化なしの場合と比較してタスク失敗リスクが90%低減されると報告されている(Indium Tech, 2026
  • コンテキスト効率: 全履歴注入と比較して、チェックポイントベースの状態管理はトークン使用量を最大90%削減できるとされている(Fastio, 2026
  • 時間的推論の向上: Statefulアーキテクチャにより、temporal reasoning(時間的推論)精度が47%、未知状況でのタスク完了率が38%改善されたという報告がある(同上)

ただし、これらの数値は特定のベンチマーク環境での測定値であり、実際の改善率はユースケースやワークロードの特性に依存します。

MCPのセッション管理と状態永続化の全体像を理解する

なぜStateful MCPサーバーが必要なのか

社内データ分析エージェントでは、ユーザーが「先月の売上データを部門別に集計して」と依頼すると、エージェントは以下のような複数ステップの処理を実行します。

  1. SQLクエリの生成と実行(数秒〜数分)
  2. 結果の集計・加工(数秒)
  3. 可視化チャートの生成(数秒)
  4. 追加分析の提案(数秒)

この一連の処理が完了する前にネットワーク断やサーバー再起動が発生した場合、Statelessなサーバーではすべてを最初からやり直す必要があります。MCPのStateful設計は、この問題をセッション管理Tasks拡張の2つの仕組みで解決します。

MCP仕様におけるセッション管理の進化

MCPのセッション管理は、仕様のバージョンアップに伴い大きく進化しています。

仕様バージョン セッション管理方式 特徴
2025-03-26 Mcp-Session-Idヘッダー導入 Streamable HTTPでのセッション識別
2025-11-25 Tasks拡張(実験的) 非同期タスク管理、ポーリングベースの状態取得
2026-07-28 RC Statelessプロトコルコア セッション依存からタスクハンドル方式へ移行

2025-11-25で導入されたTasks拡張は、ツール実行を非同期タスクとして管理し、5つの状態(working / input_required / completed / failed / cancelled)を持つステートマシンとして定義されています(MCP公式仕様)。

さらに2026-07-28のリリース候補では、プロトコルレベルのセッション自体を廃止し、Statelessなプロトコルコアへと移行する方針が示されています(MCP公式ブログ)。サーバーが状態を保持する必要がある場合は、ツールから明示的なハンドルを発行し、モデルが後続の呼び出しで引数として渡す設計が推奨されています。

3層状態管理アーキテクチャ

本記事で採用するアーキテクチャは、状態の性質に応じて3つの層に分離します。

ストレージ 保持対象 TTL 用途
L1 インメモリ(dict) 現在の実行コンテキスト リクエスト中 アクティブなツール実行の一時状態
L2 SQLite タスクチェックポイント 24時間 再接続時の状態復元
L3 Redis(オプション) 分散セッション情報 設定可能 水平スケール時のセッション共有

なぜこの構成にしたか:

  • L1のみでは再起動時にすべての状態が失われる
  • L2(SQLite)は単一インスタンスで十分な場合に最もシンプルで、ACID準拠の永続化を提供する
  • L3(Redis)は複数インスタンスでのセッション共有が必要な場合に追加する

制約条件:

SQLiteは単一プロセスからの書き込みに最適化されているため、複数プロセスから同時にチェックポイントを書き込む構成では、WAL(Write-Ahead Logging)モードの有効化が必須です。水平スケールが必要な場合はPostgreSQLまたはRedisへの移行を検討してください。

Python + FastMCPでStatefulサーバーを実装する

プロジェクト構成

まず、社内データ分析エージェント用MCPサーバーの全体構成を設計します。

data-analysis-mcp/
├── pyproject.toml
├── src/
│   └── data_analysis_mcp/
│       ├── __init__.py
│       ├── server.py          # MCPサーバー本体
│       ├── state.py           # 状態管理(チェックポイント)
│       ├── tools/
│       │   ├── __init__.py
│       │   ├── query.py       # SQLクエリ実行ツール
│       │   └── analysis.py    # 分析・集計ツール
│       └── db.py              # DB接続管理
└── tests/
    ├── test_state.py
    └── test_tools.py

Lifespan APIでDB接続を管理する

FastMCPのlifespan APIを使い、サーバー起動時にDB接続を初期化し、シャットダウン時に解放します。この設計により、各ツールはContextオブジェクト経由でDB接続にアクセスできます。

# src/data_analysis_mcp/server.py
from contextlib import asynccontextmanager
from collections.abc import AsyncIterator
from dataclasses import dataclass

import aiosqlite
from mcp.server.fastmcp import FastMCP

from data_analysis_mcp.state import CheckpointStore


@dataclass
class AppContext:
    """サーバー全体で共有するリソース"""
    checkpoint_store: CheckpointStore
    analysis_db_path: str


@asynccontextmanager
async def app_lifespan(server: FastMCP) -> AsyncIterator[AppContext]:
    checkpoint_store = CheckpointStore("checkpoints.db")
    await checkpoint_store.initialize()
    try:
        yield AppContext(
            checkpoint_store=checkpoint_store,
            analysis_db_path="sqlite:///data/analysis.db",
        )
    finally:
        await checkpoint_store.close()


mcp = FastMCP(
    "DataAnalysisMCP",
    lifespan=app_lifespan,
)

チェックポイントストアを実装する

ツール実行の途中経過をSQLiteに保存し、再接続時に復元できるようにします。各チェックポイントにはタスクIDステップ番号中間結果タイムスタンプを保持します。

# src/data_analysis_mcp/state.py
from __future__ import annotations

import json
from datetime import datetime, timezone
from enum import Enum

import aiosqlite


class TaskStatus(str, Enum):
    WORKING = "working"
    INPUT_REQUIRED = "input_required"
    COMPLETED = "completed"
    FAILED = "failed"
    CANCELLED = "cancelled"


class CheckpointStore:
    """SQLiteベースのチェックポイント永続化ストア"""

    def __init__(self, db_path: str) -> None:
        self._db_path = db_path
        self._conn: aiosqlite.Connection | None = None

    async def initialize(self) -> None:
        self._conn = await aiosqlite.connect(self._db_path)
        await self._conn.execute("PRAGMA journal_mode=WAL")
        await self._conn.execute("PRAGMA busy_timeout=5000")
        await self._conn.execute("""
            CREATE TABLE IF NOT EXISTS checkpoints (
                task_id TEXT NOT NULL,
                step INTEGER NOT NULL,
                status TEXT NOT NULL DEFAULT 'working',
                intermediate_result TEXT,
                created_at TEXT NOT NULL,
                PRIMARY KEY (task_id, step)
            )
        """)
        await self._conn.execute("""
            CREATE INDEX IF NOT EXISTS idx_checkpoints_task
            ON checkpoints (task_id, step DESC)
        """)
        await self._conn.commit()

    async def save_checkpoint(
        self,
        task_id: str,
        step: int,
        status: TaskStatus,
        intermediate_result: dict | None = None,
    ) -> None:
        assert self._conn is not None
        await self._conn.execute(
            """
            INSERT OR REPLACE INTO checkpoints
                (task_id, step, status, intermediate_result, created_at)
            VALUES (?, ?, ?, ?, ?)
            """,
            (
                task_id,
                step,
                status.value,
                json.dumps(intermediate_result) if intermediate_result else None,
                datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M:%S"),
            ),
        )
        await self._conn.commit()

    async def get_latest_checkpoint(
        self, task_id: str
    ) -> dict | None:
        assert self._conn is not None
        cursor = await self._conn.execute(
            """
            SELECT task_id, step, status, intermediate_result, created_at
            FROM checkpoints
            WHERE task_id = ?
            ORDER BY step DESC
            LIMIT 1
            """,
            (task_id,),
        )
        row = await cursor.fetchone()
        if row is None:
            return None
        return {
            "task_id": row[0],
            "step": row[1],
            "status": row[2],
            "intermediate_result": json.loads(row[3]) if row[3] else None,
            "created_at": row[4],
        }

    async def cleanup_expired(self, ttl_hours: int = 24) -> int:
        assert self._conn is not None
        cursor = await self._conn.execute(
            """
            DELETE FROM checkpoints
            WHERE created_at < datetime('now', ? || ' hours')
            """,
            (f"-{ttl_hours}",),
        )
        await self._conn.commit()
        return cursor.rowcount

    async def close(self) -> None:
        if self._conn:
            await self._conn.close()

ハマりポイント:

PRAGMA journal_mode=WALを設定しないと、読み取りと書き込みが競合した場合にdatabase is lockedエラーが発生します。WALモードにすると読み取りと書き込みが並行して実行できるようになりますが、WALファイルが肥大化する可能性があるため、定期的なPRAGMA wal_checkpoint(TRUNCATE)の実行を推奨します。

データ分析ツールにチェックポイントを組み込む

実際のSQLクエリ実行ツールに、チェックポイント保存と復元のロジックを組み込みます。長時間かかるクエリでも途中経過が保存されるため、中断後に再開できます。

# src/data_analysis_mcp/tools/query.py
from __future__ import annotations

import json
import uuid
from typing import Any

import aiosqlite
from mcp.server.fastmcp import Context

from data_analysis_mcp.server import AppContext, mcp
from data_analysis_mcp.state import CheckpointStore, TaskStatus


async def _execute_query(db_path: str, sql: str) -> list[dict[str, Any]]:
    """社内DBに対してSQLを実行し、結果を返す"""
    async with aiosqlite.connect(db_path) as db:
        db.row_factory = aiosqlite.Row
        cursor = await db.execute(sql)
        rows = await cursor.fetchall()
        columns = [desc[0] for desc in cursor.description] if cursor.description else []
        return [dict(zip(columns, row)) for row in rows]


@mcp.tool()
async def run_analysis_query(
    sql: str,
    description: str,
    task_id: str | None = None,
    ctx: Context | None = None,
) -> str:
    """社内データベースに対してSQLクエリを実行し、結果を返します。

    Args:
        sql: 実行するSQLクエリ(SELECT文のみ許可)
        description: クエリの目的の説明
        task_id: 再開時に使用するタスクID(省略時は新規生成)
    """
    app_ctx: AppContext = ctx.request_context.lifespan_context
    store: CheckpointStore = app_ctx.checkpoint_store

    if not sql.strip().upper().startswith("SELECT"):
        return "エラー: SELECT文のみ実行可能です。データ変更はできません。"

    current_task_id = task_id or str(uuid.uuid4())

    # 既存チェックポイントがあれば復元を試みる
    if task_id:
        checkpoint = await store.get_latest_checkpoint(task_id)
        if checkpoint and checkpoint["status"] == TaskStatus.COMPLETED.value:
            return (
                f"タスク {task_id} は既に完了しています。\n"
                f"結果: {checkpoint['intermediate_result']}"
            )
        if checkpoint and checkpoint["step"] >= 1:
            return (
                f"タスク {task_id} のチェックポイント(step={checkpoint['step']})から"
                f"復元しました。\n前回の中間結果: {checkpoint['intermediate_result']}"
            )

    # Step 1: クエリ検証
    await store.save_checkpoint(
        task_id=current_task_id,
        step=0,
        status=TaskStatus.WORKING,
        intermediate_result={"phase": "validation", "sql": sql},
    )

    # Step 2: クエリ実行
    try:
        results = await _execute_query(app_ctx.analysis_db_path, sql)
    except Exception as e:
        await store.save_checkpoint(
            task_id=current_task_id,
            step=1,
            status=TaskStatus.FAILED,
            intermediate_result={"error": str(e)},
        )
        return f"クエリ実行エラー (task_id={current_task_id}): {e}"

    # Step 3: 結果保存
    summary = {
        "row_count": len(results),
        "columns": list(results[0].keys()) if results else [],
        "preview": results[:10],
    }
    await store.save_checkpoint(
        task_id=current_task_id,
        step=2,
        status=TaskStatus.COMPLETED,
        intermediate_result=summary,
    )

    return json.dumps(
        {
            "task_id": current_task_id,
            "description": description,
            "row_count": len(results),
            "results": results[:100],
        },
        ensure_ascii=False,
        indent=2,
    )

再接続設計とMcp-Session-Id管理を実装する

Streamable HTTPでのセッション管理

MCPのStreamable HTTPトランスポートでは、Mcp-Session-Idヘッダーでクライアントとサーバーのセッションを紐づけます。サーバー再起動やネットワーク断の後、クライアントが古いセッションIDで接続してきた場合の処理が重要です。

再接続時のタスク復元フロー

クライアントが再接続した際に、過去のタスクIDを指定して途中から再開するパターンを実装します。ポイントは2つの復元ツールを用意することです。

# src/data_analysis_mcp/tools/analysis.py
from __future__ import annotations
import json
from mcp.server.fastmcp import Context
from data_analysis_mcp.server import AppContext, mcp


@mcp.tool()
async def resume_task(task_id: str, ctx: Context | None = None) -> str:
    """中断されたタスクを再開します。"""
    app_ctx: AppContext = ctx.request_context.lifespan_context
    store = app_ctx.checkpoint_store

    checkpoint = await store.get_latest_checkpoint(task_id)
    if checkpoint is None:
        return f"タスク {task_id} のチェックポイントが見つかりません。"

    return json.dumps(
        {
            "task_id": task_id,
            "status": checkpoint["status"],
            "last_step": checkpoint["step"],
            "intermediate_result": checkpoint["intermediate_result"],
        },
        ensure_ascii=False,
    )


@mcp.tool()
async def list_active_tasks(ctx: Context | None = None) -> str:
    """現在アクティブなタスクの一覧を返します。"""
    app_ctx: AppContext = ctx.request_context.lifespan_context
    store = app_ctx.checkpoint_store

    assert store._conn is not None
    cursor = await store._conn.execute("""
        SELECT task_id, MAX(step) as last_step, status, created_at
        FROM checkpoints
        GROUP BY task_id
        HAVING status NOT IN ('completed', 'failed', 'cancelled')
        ORDER BY created_at DESC LIMIT 20
    """)
    rows = await cursor.fetchall()
    tasks = [{"task_id": r[0], "last_step": r[1], "status": r[2], "created_at": r[3]} for r in rows]
    return json.dumps(tasks, ensure_ascii=False, indent=2)

MCP Tasks拡張で長時間実行ツールを管理する

Tasks拡張の仕組み

MCP 2025-11-25で導入されたTasks拡張は、ツール呼び出しを非同期タスクとして扱います(MCP公式仕様)。通常のツール呼び出しがリクエスト→レスポンスの同期パターンであるのに対し、Tasks拡張ではリクエスト→タスクID返却→ポーリング→結果取得という非同期パターンになります。

Tasks対応ツールの実装パターン

FastMCPでTasks拡張に対応したツールを実装する場合、execution.taskSupport"optional"に設定します。クライアントがタスクモードで呼び出した場合は非同期で実行し、通常の呼び出しでは同期で応答します。

# src/data_analysis_mcp/tools/query.py (Tasks拡張対応版)
import asyncio
from typing import Any


# タスクストア(インメモリ + SQLiteチェックポイント)
_running_tasks: dict[str, asyncio.Task] = {}


async def _execute_long_query(
    db_path: str,
    sql: str,
    task_id: str,
    store: CheckpointStore,
) -> dict[str, Any]:
    """長時間クエリをバックグラウンドで実行する"""
    await store.save_checkpoint(
        task_id=task_id,
        step=0,
        status=TaskStatus.WORKING,
        intermediate_result={"phase": "executing", "sql": sql},
    )

    try:
        results = await _execute_query(db_path, sql)
        summary = {
            "row_count": len(results),
            "columns": list(results[0].keys()) if results else [],
            "data": results[:100],
        }
        await store.save_checkpoint(
            task_id=task_id,
            step=1,
            status=TaskStatus.COMPLETED,
            intermediate_result=summary,
        )
        return summary
    except Exception as e:
        await store.save_checkpoint(
            task_id=task_id,
            step=1,
            status=TaskStatus.FAILED,
            intermediate_result={"error": str(e)},
        )
        raise

タスクのTTLとクリーンアップ

MCPのTasks仕様では、各タスクに**TTL(Time To Live)**を設定し、有効期限が切れたタスクはサーバーが自動的に削除できると定義されています。これをチェックポイントストアのクリーンアップと連携させます。

# src/data_analysis_mcp/cleanup.py
import asyncio

from data_analysis_mcp.state import CheckpointStore


async def periodic_cleanup(store: CheckpointStore, interval_hours: int = 1) -> None:
    """期限切れチェックポイントを定期的に削除する"""
    while True:
        deleted = await store.cleanup_expired(ttl_hours=24)
        if deleted > 0:
            print(f"Cleaned up {deleted} expired checkpoints")
        await asyncio.sleep(interval_hours * 3600)

よくある問題と解決方法

社内データ分析エージェントをStateful MCPサーバーで運用する際に遭遇しやすい問題と、その対処法をまとめます。

問題 原因 解決方法
Session not foundエラーが頻発 サーバー再起動後に古いMcp-Session-Idでリクエスト クライアント側で404応答時にinitializeを再送信し、新セッションを確立する
チェックポイントDBが肥大化 TTLクリーンアップが未設定 cleanup_expired()を定期実行(推奨: 1時間ごと)
database is lockedエラー SQLiteのデフォルトジャーナルモードでの並行アクセス PRAGMA journal_mode=WALPRAGMA busy_timeout=5000を設定
再接続後にタスクIDが不明 クライアントがタスクIDを保持していない list_active_tasksツールで未完了タスクを照会可能にする
水平スケール時にセッション不整合 インメモリセッションが各インスタンスに分散 L3層(Redis)にセッション情報を移動するか、ロードバランサーでsticky sessionを設定する
Tasks対応クライアントが少ない Tasks拡張が実験的ステータス(2025-11-25時点) taskSupport: "optional"に設定し、同期・非同期の両方に対応する

最初はtaskSupport: "required"で運用したところ、一部のMCPクライアント(VS Code拡張の古いバージョン等)が非対応でツール呼び出しに失敗するケースがありました。 本番環境では"optional"に設定し、クライアントの対応状況に応じてフォールバックするのが安全です。

まとめと次のステップ

まとめ:

  • MCPのセッション管理Mcp-Session-Id)とTasks拡張(5状態のステートマシン)を組み合わせることで、社内データ分析エージェントの状態永続化と再接続を実現できる
  • 状態管理は3層アーキテクチャ(インメモリ / SQLite / Redis)で設計し、スケーラビリティ要件に応じて層を追加する
  • 2026-07-28のMCP仕様RCではStatelessプロトコルコアへ移行する方向性が示されており、タスクハンドル方式への移行が推奨されている
  • チェックポイントのTTL管理WALモード設定は、本番運用での安定性に直結する
  • taskSupport"optional"に設定して、非対応クライアントへのフォールバックを確保する

次にやるべきこと:

  • MCP Tasks拡張の公式仕様(Tasks - Model Context Protocol)を通読し、input_required状態を活用したインタラクティブな分析フローを設計する
  • FastMCPのlifespan APIでRedis接続を追加し、セッション情報の分散共有を検証する
  • 2026-07-28のMCP仕様RC(公式ブログ)を確認し、Stateless化への移行計画を立てる

参考


関連する深掘り記事

この記事で紹介した技術について、さらに深掘りした記事を書きました:

GitHubで編集を提案

Discussion