MCP资源模板实战:如何用Python动态生成用户信息(附避坑指南)

最近在构建AI应用时,我遇到了一个常见但棘手的问题:如何让大语言模型访问动态变化的数据?比如用户信息、实时库存、个性化配置等。传统的静态资源方案显然不够用,每次数据变化都需要重新部署服务,这在实际项目中几乎不可行。

经过几轮探索,我发现了MCP(Model Context Protocol)中的Resource Template机制,它完美解决了这个问题。Resource Template允许你定义带参数的URI模板,客户端可以根据不同参数请求不同的资源内容。这就像为你的AI应用装上了“动态数据引擎”,让大语言模型能够实时访问个性化的上下文信息。

今天我就来分享如何用Python构建一个完整的Resource Template服务,从基础概念到实战代码,再到那些容易踩坑的细节,帮你快速掌握这项实用技术。

1. 理解Resource Template的核心价值

Resource Template是MCP协议中一个相对高级但极其实用的功能。与静态Resource不同,它允许你在URI中嵌入参数,根据客户端传入的不同参数值返回不同的内容。这种设计模式在现代AI应用中有着广泛的应用场景。

1.1 为什么需要动态资源?

在真实业务场景中,数据很少是静态不变的。考虑以下几个典型需求:

  • 用户个性化数据:每个用户都有自己的配置、历史记录、偏好设置
  • 实时业务数据:库存状态、订单信息、交易记录等随时间变化
  • 参数化查询:根据不同的查询条件返回不同的数据集
  • 多租户系统:不同客户看到不同的数据视图

如果为每个可能的参数组合都创建一个静态资源,那将是灾难性的。Resource Template通过URI模板解决了这个问题,让单个资源定义能够服务无限种参数组合。

1.2 Resource Template的工作原理

Resource Template的核心思想很简单:定义URI模板 → 客户端传入参数 → 服务端动态生成内容。整个过程可以分为三个关键步骤:

  1. 模板定义:服务端定义包含参数的URI模板,如 user://{user_id}/profile
  2. 参数传递:客户端在请求时提供具体的参数值,如 user_id="123"
  3. 动态生成:服务端根据参数值生成对应的资源内容

这种模式与Web开发中的RESTful API路由参数非常相似,但专门为AI应用场景优化。MCP协议标准化了这一交互过程,确保不同客户端和服务端之间的兼容性。

1.3 与静态Resource的对比

为了更清晰地理解Resource Template的优势,我们来看一个对比表格:

特性 静态Resource Resource Template
URI格式 固定字符串 包含参数的模板字符串
数据源 单一、静态 动态、参数驱动
适用场景 配置文件、文档、静态数据 用户数据、实时信息、个性化内容
扩展性 有限,需要预定义所有资源 高,通过参数支持无限组合
维护成本 数据变化时需要更新资源定义 只需维护生成逻辑,数据变化自动适应

从实际开发经验来看,Resource Template特别适合那些需要根据用户身份、查询条件或时间因素返回不同数据的场景。比如用户管理系统、电商商品目录、数据分析报告等。

2. 构建Resource Template服务端

现在让我们进入实战环节。我将通过一个完整的用户信息查询服务来演示Resource Template的实现过程。这个服务将根据用户ID返回对应的用户资料,包括基本信息、权限设置和个性化配置。

2.1 环境准备与项目结构

首先确保你的开发环境已经准备就绪。我推荐使用Python 3.10+版本,因为MCP的异步特性在较新版本中支持更好。

# 创建项目目录
mkdir user-info-mcp
cd user-info-mcp

# 创建虚拟环境
python -m venv venv

# 激活虚拟环境(Linux/Mac)
source venv/bin/activate

# 激活虚拟环境(Windows)
venv\Scripts\activate

# 安装核心依赖
pip install mcp fastmcp aiofiles httpx pydantic

项目结构应该清晰合理,便于维护和扩展:

user-info-mcp/
├── server.py          # 主服务文件
├── client.py          # 客户端示例
├── models.py          # 数据模型定义
├── data/
│   └── users.json     # 模拟用户数据
└── requirements.txt   # 依赖列表

models.py中,我们先定义用户数据模型,这有助于保持代码的类型安全和可维护性:

from pydantic import BaseModel, Field
from typing import Optional, Dict, Any
from datetime import datetime

class UserProfile(BaseModel):
    """用户基本信息模型"""
    user_id: str = Field(..., description="用户唯一标识")
    username: str = Field(..., description="用户名")
    email: str = Field(..., description="邮箱地址")
    created_at: datetime = Field(default_factory=datetime.now, description="创建时间")
    last_login: Optional[datetime] = Field(None, description="最后登录时间")
    
class UserPreferences(BaseModel):
    """用户偏好设置模型"""
    theme: str = Field("light", description="界面主题")
    language: str = Field("zh-CN", description="界面语言")
    notifications_enabled: bool = Field(True, description="是否启用通知")
    
class UserData(BaseModel):
    """完整用户数据模型"""
    profile: UserProfile
    preferences: UserPreferences
    metadata: Dict[str, Any] = Field(default_factory=dict, description="扩展元数据")

2.2 实现Resource Template服务

接下来是核心部分——实现Resource Template服务。我将创建一个完整的服务端,支持多种用户数据查询模式。

# server.py
import json
import asyncio
from typing import Dict, Any
from pathlib import Path
from mcp.server.fastmcp import FastMCP
from models import UserData, UserProfile, UserPreferences

# 初始化FastMCP应用
app = FastMCP("user-info-service")

# 模拟用户数据存储
class UserStore:
    """用户数据存储管理器"""
    
    def __init__(self, data_file: str = "data/users.json"):
        self.data_file = Path(data_file)
        self._users: Dict[str, UserData] = {}
        self._load_data()
    
    def _load_data(self):
        """从文件加载用户数据"""
        if self.data_file.exists():
            with open(self.data_file, 'r', encoding='utf-8') as f:
                raw_data = json.load(f)
                for user_id, user_dict in raw_data.items():
                    self._users[user_id] = UserData(**user_dict)
        else:
            # 创建示例数据
            self._create_sample_data()
    
    def _create_sample_data(self):
        """创建示例用户数据"""
        sample_users = {
            "1001": UserData(
                profile=UserProfile(
                    user_id="1001",
                    username="张三",
                    email="zhangsan@example.com"
                ),
                preferences=UserPreferences(
                    theme="dark",
                    language="zh-CN",
                    notifications_enabled=True
                ),
                metadata={"role": "admin", "department": "技术部"}
            ),
            "1002": UserData(
                profile=UserProfile(
                    user_id="1002",
                    username="李四",
                    email="lisi@example.com"
                ),
                preferences=UserPreferences(
                    theme="light",
                    language="en-US",
                    notifications_enabled=False
                ),
                metadata={"role": "user", "department": "市场部"}
            )
        }
        self._users = sample_users
        self._save_data()
    
    def _save_data(self):
        """保存数据到文件"""
        self.data_file.parent.mkdir(exist_ok=True)
        with open(self.data_file, 'w', encoding='utf-8') as f:
            # 使用Pydantic的dict方法确保序列化正确
            data_dict = {
                uid: user_data.dict() 
                for uid, user_data in self._users.items()
            }
            json.dump(data_dict, f, ensure_ascii=False, indent=2, default=str)
    
    def get_user(self, user_id: str) -> UserData:
        """根据用户ID获取用户数据"""
        if user_id not in self._users:
            raise ValueError(f"用户 {user_id} 不存在")
        return self._users[user_id]
    
    def search_users(self, **filters) -> Dict[str, UserData]:
        """根据条件搜索用户"""
        results = {}
        for uid, user in self._users.items():
            match = True
            for key, value in filters.items():
                # 支持嵌套属性查询
                if '.' in key:
                    parts = key.split('.')
                    obj = user
                    for part in parts:
                        if hasattr(obj, part):
                            obj = getattr(obj, part)
                        else:
                            match = False
                            break
                    if match and str(obj) != str(value):
                        match = False
                elif hasattr(user, key):
                    if getattr(user, key) != value:
                        match = False
                else:
                    match = False
                
                if not match:
                    break
            
            if match:
                results[uid] = user
        
        return results

# 初始化用户存储
user_store = UserStore()

# 定义Resource Template - 基础用户信息
@app.resource(
    uri="user://{user_id}",
    name="user_basic_info",
    description="根据用户ID获取基础用户信息\n:param user_id: 用户的唯一标识符",
    mime_type="application/json"
)
async def get_user_basic_info(user_id: str) -> Dict[str, Any]:
    """
    获取用户基础信息
    
    这个Resource Template根据用户ID返回用户的基本资料,
    包括用户名、邮箱和创建时间等核心信息。
    """
    try:
        user_data = user_store.get_user(user_id)
        return {
            "user_id": user_data.profile.user_id,
            "username": user_data.profile.username,
            "email": user_data.profile.email,
            "created_at": user_data.profile.created_at.isoformat(),
            "last_login": user_data.profile.last_login.isoformat() if user_data.profile.last_login else None
        }
    except ValueError as e:
        return {"error": str(e), "user_id": user_id}

# 定义Resource Template - 完整用户档案
@app.resource(
    uri="user://{user_id}/profile",
    name="user_full_profile",
    description="获取用户的完整档案信息\n:param user_id: 用户的唯一标识符",
    mime_type="application/json"
)
async def get_user_full_profile(user_id: str) -> Dict[str, Any]:
    """
    获取用户完整档案
    
    返回用户的全部信息,包括个人资料、偏好设置和元数据。
    适合需要完整用户上下文的场景。
    """
    try:
        user_data = user_store.get_user(user_id)
        return user_data.dict()
    except ValueError as e:
        return {"error": str(e), "user_id": user_id}

# 定义Resource Template - 用户偏好设置
@app.resource(
    uri="user://{user_id}/preferences",
    name="user_preferences",
    description="获取用户的个性化偏好设置\n:param user_id: 用户的唯一标识符",
    mime_type="application/json"
)
async def get_user_preferences(user_id: str) -> Dict[str, Any]:
    """
    获取用户偏好设置
    
    专门返回用户的界面和交互偏好,
    适用于个性化UI配置的场景。
    """
    try:
        user_data = user_store.get_user(user_id)
        return user_data.preferences.dict()
    except ValueError as e:
        return {"error": str(e), "user_id": user_id}

# 定义Resource Template - 带过滤条件的用户搜索
@app.resource(
    uri="users://search?department={department}&role={role}",
    name="user_search",
    description="根据部门和角色筛选用户\n:param department: 部门名称\n:param role: 用户角色",
    mime_type="application/json"
)
async def search_users_by_filters(department: str = None, role: str = None) -> Dict[str, Any]:
    """
    多条件用户搜索
    
    支持根据部门和角色组合筛选用户,
    参数都是可选的,可以单独使用或组合使用。
    """
    filters = {}
    if department:
        filters["metadata.department"] = department
    if role:
        filters["metadata.role"] = role
    
    results = user_store.search_users(**filters)
    
    return {
        "filters": filters,
        "count": len(results),
        "users": {
            uid: user_data.dict()
            for uid, user_data in results.items()
        }
    }

if __name__ == "__main__":
    # 启动SSE服务,便于调试
    app.run(transport="sse", host="127.0.0.1", port=8000)

这个服务端实现了几个关键功能:

  1. 模块化设计:使用Pydantic模型确保数据结构的类型安全
  2. 错误处理:对不存在的用户ID返回明确的错误信息
  3. 多种资源模板:提供不同粒度的用户信息访问
  4. 参数化搜索:支持多条件组合查询

注意:在实际生产环境中,你应该将用户数据存储在数据库或外部服务中,这里使用JSON文件只是为了演示方便。

2.3 服务端的高级特性

为了让Resource Template服务更加健壮和实用,我通常会添加一些高级特性:

# 在server.py中添加以下内容
import time
from functools import lru_cache
from mcp.server.fastmcp import ResourceTemplate

# 添加缓存机制
@lru_cache(maxsize=128)
def get_user_cached(user_id: str) -> Dict[str, Any]:
    """带缓存的用户数据获取"""
    # 模拟数据库查询延迟
    time.sleep(0.1)
    user_data = user_store.get_user(user_id)
    return user_data.dict()

# 带缓存的Resource Template
@app.resource(
    uri="user://{user_id}/cached",
    name="user_cached_info",
    description="带缓存的用户信息查询\n:param user_id: 用户的唯一标识符",
    mime_type="application/json"
)
async def get_user_cached_info(user_id: str) -> Dict[str, Any]:
    """
    带缓存的用户信息查询
    
    使用LRU缓存提高频繁访问的性能,
    适合用户信息变化不频繁的场景。
    """
    try:
        return get_user_cached(user_id)
    except ValueError as e:
        return {"error": str(e), "user_id": user_id}

# 批量用户查询模板
@app.resource(
    uri="users://batch?ids={user_ids}",
    name="batch_user_info",
    description="批量获取多个用户的信息\n:param user_ids: 用户ID列表,用逗号分隔",
    mime_type="application/json"
)
async def get_batch_user_info(user_ids: str) -> Dict[str, Any]:
    """
    批量用户信息查询
    
    支持一次查询多个用户的信息,
    减少网络请求次数,提高效率。
    """
    ids = [uid.strip() for uid in user_ids.split(',')]
    results = {}
    errors = []
    
    for user_id in ids:
        try:
            user_data = user_store.get_user(user_id)
            results[user_id] = user_data.dict()
        except ValueError as e:
            errors.append({"user_id": user_id, "error": str(e)})
    
    return {
        "success_count": len(results),
        "error_count": len(errors),
        "results": results,
        "errors": errors
    }

这些高级特性在实际项目中非常有用:

  • 缓存机制:对于不经常变化的数据,缓存可以显著提高性能
  • 批量查询:减少客户端请求次数,优化网络传输
  • 错误隔离:批量操作中部分失败不影响整体结果

3. 客户端实现与参数处理

服务端准备好了,接下来看看客户端如何调用这些Resource Template。这是最容易出问题的环节,特别是参数传递和URI构建。

3.1 基础客户端实现

首先创建一个基础的MCP客户端,用于连接我们的Resource Template服务:

# client.py
import asyncio
import json
from typing import Dict, Any, Optional
from mcp import ClientSession
from mcp.client.sse import sse_client
from openai import OpenAI
from contextlib import AsyncExitStack

class UserInfoClient:
    """用户信息MCP客户端"""
    
    def __init__(self, api_key: str, base_url: str = "https://api.deepseek.com"):
        """
        初始化MCP客户端
        
        Args:
            api_key: DeepSeek API密钥
            base_url: API基础URL
        """
        self.deepseek = OpenAI(
            api_key=api_key,
            base_url=base_url
        )
        self.exit_stack = AsyncExitStack()
        self.session: Optional[ClientSession] = None
        self.resource_templates: Dict[str, Dict[str, Any]] = {}
    
    async def connect(self, server_url: str = "http://127.0.0.1:8000/sse"):
        """
        连接到MCP服务器
        
        Args:
            server_url: SSE服务器URL
        """
        try:
            # 创建SSE连接
            read_stream, write_stream = await self.exit_stack.enter_async_context(
                sse_client(server_url)
            )
            
            # 创建会话
            self.session = await self.exit_stack.enter_async_context(
                ClientSession(read_stream, write_stream)
            )
            
            # 初始化会话
            await self.session.initialize()
            print("✅ 成功连接到MCP服务器")
            
            # 发现可用的Resource Template
            await self._discover_resources()
            
        except Exception as e:
            print(f"❌ 连接失败: {e}")
            raise
    
    async def _discover_resources(self):
        """发现服务器提供的Resource Template"""
        if not self.session:
            raise RuntimeError("会话未初始化")
        
        # 获取Resource Template列表
        response = await self.session.list_resource_templates()
        
        for template in response.resourceTemplates:
            self.resource_templates[template.name] = {
                "uri_template": template.uriTemplate,
                "description": template.description,
                "mime_type": template.mimeType
            }
        
        print(f"📋 发现 {len(self.resource_templates)} 个Resource Template:")
        for name, info in self.resource_templates.items():
            print(f"  • {name}: {info['description']}")
    
    async def query_user_info(self, user_query: str) -> str:
        """
        查询用户信息
        
        Args:
            user_query: 用户查询的自然语言描述
            
        Returns:
            大模型生成的回答
        """
        if not self.session:
            raise RuntimeError("请先调用connect()方法连接服务器")
        
        # 准备Function Calling格式的工具描述
        functions = []
        for name, info in self.resource_templates.items():
            functions.append({
                "type": "function",
                "function": {
                    "name": name,
                    "description": info["description"],
                    "input_schema": None  # Resource Template不需要输入模式
                }
            })
        
        # 构建消息
        messages = [{
            "role": "system",
            "content": "你是一个用户信息查询助手。你可以使用可用的资源模板来获取用户信息。"
        }, {
            "role": "user",
            "content": user_query
        }]
        
        # 调用大模型
        response = self.deepseek.chat.completions.create(
            model="deepseek-chat",
            messages=messages,
            tools=functions,
            tool_choice="auto"
        )
        
        # 处理响应
        choice = response.choices[0]
        
        if choice.finish_reason == "tool_calls":
            tool_calls = choice.message.tool_calls
            
            for tool_call in tool_calls:
                function = tool_call.function
                function_name = function.name
                function_args = json.loads(function.arguments) if function.arguments else {}
                
                print(f"🔧 调用工具: {function_name}")
                print(f"   参数: {function_args}")
                
                # 获取对应的URI模板
                uri_template = self.resource_templates[function_name]["uri_template"]
                
                # 构建完整的URI
                try:
                    # 这里是最关键的一步:将参数填充到URI模板中
                    full_uri = uri_template.format(**function_args)
                    print(f"   URI: {full_uri}")
                    
                    # 读取资源
                    resource_response = await self.session.read_resource(full_uri)
                    
                    # 提取资源内容
                    if resource_response.contents:
                        content = resource_response.contents[0]
                        if hasattr(content, 'text'):
                            result = content.text
                        elif hasattr(content, 'blob'):
                            # 处理二进制内容
                            import base64
                            result = base64.b64decode(content.blob).decode('utf-8')
                        else:
                            result = str(content)
                        
                        # 将结果返回给大模型
                        messages.append(choice.message.model_dump())
                        messages.append({
                            "role": "tool",
                            "content": result,
                            "tool_call_id": tool_call.id
                        })
                        
                        # 获取最终回答
                        final_response = self.deepseek.chat.completions.create(
                            model="deepseek-chat",
                            messages=messages
                        )
                        
                        return final_response.choices[0].message.content
                
                except KeyError as e:
                    error_msg = f"参数错误: 缺少必要的参数 {e}"
                    return error_msg
                except Exception as e:
                    error_msg = f"资源读取失败: {str(e)}"
                    return error_msg
        
        # 如果没有工具调用,直接返回大模型的回答
        return choice.message.content
    
    async def close(self):
        """关闭连接"""
        await self.exit_stack.aclose()
        print("👋 连接已关闭")

# 使用示例
async def main():
    # 初始化客户端
    client = UserInfoClient(
        api_key="your-deepseek-api-key"  # 替换为你的API密钥
    )
    
    try:
        # 连接服务器
        await client.connect()
        
        # 测试查询
        queries = [
            "获取用户ID为1001的信息",
            "查看用户1002的偏好设置",
            "搜索技术部的所有用户",
            "批量获取用户1001和1002的完整档案"
        ]
        
        for query in queries:
            print(f"\n📝 查询: {query}")
            print("-" * 50)
            result = await client.query_user_info(query)
            print(f"💬 回答: {result}")
            print("-" * 50)
            
    finally:
        await client.close()

if __name__ == "__main__":
    asyncio.run(main())

这个客户端实现了完整的Resource Template调用流程:

  1. 连接管理:自动处理SSE连接和会话生命周期
  2. 资源发现:动态获取服务器提供的Resource Template
  3. 参数处理:正确解析和填充URI模板参数
  4. 错误处理:优雅处理各种异常情况

3.2 参数处理的常见问题与解决方案

在实际使用中,参数处理是最容易出问题的地方。下面我总结了一些常见问题及其解决方案:

问题1:参数类型不匹配

Resource Template的参数都是字符串类型,但实际业务中可能需要其他类型。解决方案是在服务端进行类型转换:

# 在服务端添加类型转换逻辑
@app.resource(
    uri="user://{user_id}/posts?limit={limit}",
    name="user_posts",
    description="获取用户的帖子\n:param user_id: 用户ID\n:param limit: 返回数量限制",
    mime_type="application/json"
)
async def get_user_posts(user_id: str, limit: str = "10") -> Dict[str, Any]:
    """
    获取用户帖子列表
    
    支持分页参数,limit参数会自动转换为整数类型
    """
    try:
        limit_int = int(limit)
        if limit_int <= 0:
            limit_int = 10
        if limit_int > 100:
            limit_int = 100
    except ValueError:
        limit_int = 10
    
    # 使用转换后的limit_int进行查询
    # ...
问题2:可选参数处理

有些参数是可选的,但URI模板需要完整的参数。解决方案是使用查询参数格式:

# 使用查询参数支持可选参数
@app.resource(
    uri="users://search?name={name}&age={age}&city={city}",
    name="advanced_user_search",
    description="高级用户搜索\n:param name: 用户名(可选)\n:param age: 年龄(可选)\n:param city: 城市(可选)",
    mime_type="application/json"
)
async def advanced_search(name: str = None, age: str = None, city: str = None) -> Dict[str, Any]:
    """
    高级用户搜索
    
    所有参数都是可选的,可以任意组合
    """
    filters = {}
    if name:
        filters["name"] = name
    if age:
        try:
            filters["age"] = int(age)
        except ValueError:
            pass
    if city:
        filters["city"] = city
    
    # 构建查询
    # ...
问题3:特殊字符编码

URI中的特殊字符需要正确编码。MCP客户端通常会处理这个问题,但为了安全起见,可以在服务端进行验证:

from urllib.parse import unquote

@app.resource(
    uri="search://{query}",
    name="global_search",
    description="全局搜索\n:param query: 搜索关键词",
    mime_type="application/json"
)
async def global_search(query: str) -> Dict[str, Any]:
    """
    全局搜索
    
    自动处理URL编码的查询参数
    """
    # 解码URL编码的参数
    decoded_query = unquote(query)
    
    # 进行搜索
    # ...

3.3 客户端优化技巧

为了让客户端更加健壮和易用,我通常会添加以下优化:

# 在UserInfoClient类中添加以下方法
class UserInfoClient:
    # ... 之前的代码 ...
    
    async def query_with_retry(self, query: str, max_retries: int = 3) -> str:
        """
        带重试机制的查询
        
        Args:
            query: 查询语句
            max_retries: 最大重试次数
            
        Returns:
            查询结果
        """
        for attempt in range(max_retries):
            try:
                return await self.query_user_info(query)
            except Exception as e:
                if attempt == max_retries - 1:
                    raise
                print(f"⚠️ 查询失败,第{attempt + 1}次重试: {e}")
                await asyncio.sleep(1 * (attempt + 1))  # 指数退避
    
    async def batch_query(self, queries: list) -> Dict[str, str]:
        """
        批量查询
        
        Args:
            queries: 查询语句列表
            
        Returns:
            查询结果字典
        """
        results = {}
        tasks = []
        
        for query in queries:
            task = asyncio.create_task(self.query_user_info(query))
            tasks.append((query, task))
        
        for query, task in tasks:
            try:
                results[query] = await task
            except Exception as e:
                results[query] = f"查询失败: {str(e)}"
        
        return results
    
    def get_available_templates(self) -> Dict[str, str]:
        """
        获取可用的Resource Template描述
        
        Returns:
            模板名称到描述的映射
        """
        return {
            name: info["description"]
            for name, info in self.resource_templates.items()
        }

这些优化措施包括:

  • 重试机制:网络不稳定时的自动重试
  • 批量查询:并行处理多个查询提高效率
  • 模板发现:方便用户了解可用的查询功能

4. 实战案例:构建用户画像系统

现在让我们把这些知识应用到一个实际场景中:构建一个用户画像系统。这个系统能够根据用户行为数据动态生成用户画像,支持多种维度的分析。

4.1 系统架构设计

用户画像系统需要处理多种数据源和复杂的分析逻辑。我设计了以下架构:

用户画像系统架构:
1. 数据层:用户行为日志、基本信息、偏好数据
2. 服务层:Resource Template服务、分析引擎、缓存
3. 应用层:客户端应用、管理界面、API网关

4.2 实现用户画像Resource Template

# user_profile_system.py
import asyncio
from datetime import datetime, timedelta
from typing import Dict, Any, List
from mcp.server.fastmcp import FastMCP
import json

app = FastMCP("user-profile-system")

# 模拟用户行为数据
class UserBehaviorAnalyzer:
    """用户行为分析器"""
    
    def __init__(self):
        self.behavior_data = self._load_behavior_data()
    
    def _load_behavior_data(self) -> Dict[str, List[Dict]]:
        """加载模拟行为数据"""
        # 在实际项目中,这里应该连接数据库或数据仓库
        return {
            "1001": [
                {"action": "login", "timestamp": "2024-01-15T08:30:00", "duration": 300},
                {"action": "view_product", "timestamp": "2024-01-15T09:15:00", "product_id": "P001"},
                {"action": "purchase", "timestamp": "2024-01-15T10:00:00", "amount": 299.99},
                {"action": "share", "timestamp": "2024-01-15T14:30:00", "platform": "wechat"},
            ],
            "1002": [
                {"action": "login", "timestamp": "2024-01-15T09:00:00", "duration": 180},
                {"action": "view_article", "timestamp": "2024-01-15T09:30:00", "article_id": "A001"},
                {"action": "comment", "timestamp": "2024-01-15T10:15:00", "content": "很好的文章"},
                {"action": "logout", "timestamp": "2024-01-15T11:00:00", "duration": 600},
            ]
        }
    
    def analyze_behavior_pattern(self, user_id: str, days: int = 7) -> Dict[str, Any]:
        """分析用户行为模式"""
        if user_id not in self.behavior_data:
            return {"error": f"用户 {user_id} 无行为数据"}
        
        behaviors = self.behavior_data[user_id]
        
        # 计算各种指标
        total_sessions = len([b for b in behaviors if b["action"] == "login"])
        total_purchases = len([b for b in behaviors if b["action"] == "purchase"])
        total_shares = len([b for b in behaviors if b["action"] == "share"])
        
        # 计算活跃时间段
        login_times = [
            datetime.fromisoformat(b["timestamp"]).hour
            for b in behaviors if b["action"] == "login"
        ]
        active_hours = self._calculate_active_hours(login_times)
        
        return {
            "user_id": user_id,
            "analysis_period_days": days,
            "metrics": {
                "total_sessions": total_sessions,
                "total_purchases": total_purchases,
                "total_shares": total_shares,
                "avg_session_duration": self._calculate_avg_duration(behaviors),
                "preferred_action": self._find_most_common_action(behaviors),
            },
            "patterns": {
                "active_hours": active_hours,
                "behavior_frequency": self._calculate_frequency(behaviors),
                "engagement_score": self._calculate_engagement_score(behaviors),
            },
            "recommendations": self._generate_recommendations(behaviors)
        }
    
    def _calculate_active_hours(self, hours: List[int]) -> List[int]:
        """计算活跃时间段"""
        if not hours:
            return []
        
        # 简单的活跃时间段分析
        hour_counts = {}
        for hour in hours:
            hour_counts[hour] = hour_counts.get(hour, 0) + 1
        
        # 返回出现次数最多的3个时间段
        sorted_hours = sorted(hour_counts.items(), key=lambda x: x[1], reverse=True)
        return [hour for hour, _ in sorted_hours[:3]]
    
    def _calculate_avg_duration(self, behaviors: List[Dict]) -> float:
        """计算平均会话时长"""
        durations = [b.get("duration", 0) for b in behaviors if "duration" in b]
        return sum(durations) / len(durations) if durations else 0
    
    def _find_most_common_action(self, behaviors: List[Dict]) -> str:
        """找出最常见的行为"""
        action_counts = {}
        for behavior in behaviors:
            action = behavior["action"]
            action_counts[action] = action_counts.get(action, 0) + 1
        
        if action_counts:
            return max(action_counts.items(), key=lambda x: x[1])[0]
        return "unknown"
    
    def _calculate_frequency(self, behaviors: List[Dict]) -> Dict[str, float]:
        """计算行为频率"""
        # 简化实现
        return {
            "daily": len(behaviors) / 30,  # 假设30天数据
            "weekly": len(behaviors) / 4,
            "monthly": len(behaviors)
        }
    
    def _calculate_engagement_score(self, behaviors: List[Dict]) -> float:
        """计算参与度分数"""
        # 简化评分算法
        score = 0
        for behavior in behaviors:
            if behavior["action"] == "purchase":
                score += 10
            elif behavior["action"] == "share":
                score += 5
            elif behavior["action"] == "comment":
                score += 3
            elif behavior["action"] == "login":
                score += 1
        
        return min(score / 10, 10.0)  # 归一化到0-10分
    
    def _generate_recommendations(self, behaviors: List[Dict]) -> List[str]:
        """生成个性化推荐"""
        recommendations = []
        actions = [b["action"] for b in behaviors]
        
        if "purchase" not in actions:
            recommendations.append("尝试购买商品以获得更好体验")
        
        if "share" not in actions:
            recommendations.append("分享内容可以获得积分奖励")
        
        if len([b for b in behaviors if b["action"] == "login"]) < 5:
            recommendations.append("增加登录频率可以解锁更多功能")
        
        return recommendations

# 初始化分析器
analyzer = UserBehaviorAnalyzer()

# 用户画像Resource Template
@app.resource(
    uri="profile://{user_id}/behavior?period={period_days}",
    name="user_behavior_analysis",
    description="分析用户行为模式\n:param user_id: 用户ID\n:param period_days: 分析周期(天数)",
    mime_type="application/json"
)
async def analyze_user_behavior(user_id: str, period_days: str = "7") -> Dict[str, Any]:
    """
    用户行为分析
    
    根据指定周期分析用户行为模式,
    生成详细的画像报告。
    """
    try:
        days = int(period_days)
        if days <= 0:
            days = 7
        if days > 365:
            days = 365
    except ValueError:
        days = 7
    
    analysis = analyzer.analyze_behavior_pattern(user_id, days)
    
    # 添加时间戳和元数据
    analysis["metadata"] = {
        "generated_at": datetime.now().isoformat(),
        "analysis_period": f"{days}天",
        "data_source": "模拟行为数据"
    }
    
    return analysis

# 用户标签生成
@app.resource(
    uri="profile://{user_id}/tags",
    name="user_tags",
    description="生成用户标签\n:param user_id: 用户ID",
    mime_type="application/json"
)
async def generate_user_tags(user_id: str) -> Dict[str, Any]:
    """
    生成用户标签
    
    基于用户行为数据自动生成标签,
    用于用户分群和个性化推荐。
    """
    analysis = analyzer.analyze_behavior_pattern(user_id)
    
    tags = []
    
    # 基于行为生成标签
    metrics = analysis.get("metrics", {})
    patterns = analysis.get("patterns", {})
    
    if metrics.get("total_purchases", 0) > 5:
        tags.append("高价值客户")
    
    if metrics.get("total_shares", 0) > 3:
        tags.append("活跃分享者")
    
    engagement_score = patterns.get("engagement_score", 0)
    if engagement_score > 7:
        tags.append("高参与度用户")
    elif engagement_score > 4:
        tags.append("中等参与度用户")
    else:
        tags.append("低参与度用户")
    
    # 基于活跃时间生成标签
    active_hours = patterns.get("active_hours", [])
    if any(9 <= hour <= 17 for hour in active_hours[:2]):
        tags.append("工作日活跃")
    if any(18 <= hour <= 23 or 0 <= hour <= 8 for hour in active_hours[:2]):
        tags.append("夜间活跃")
    
    return {
        "user_id": user_id,
        "tags": tags,
        "tag_count": len(tags),
        "generated_at": datetime.now().isoformat(),
        "tag_details": {
            "behavior_based": [t for t in tags if t in ["高价值客户", "活跃分享者", "高参与度用户"]],
            "time_based": [t for t in tags if t in ["工作日活跃", "夜间活跃"]]
        }
    }

# 用户相似度分析
@app.resource(
    uri="profile://compare?user1={user1_id}&user2={user2_id}",
    name="user_comparison",
    description="比较两个用户的画像\n:param user1_id: 第一个用户ID\n:param user2_id: 第二个用户ID",
    mime_type="application/json"
)
async def compare_users(user1_id: str, user2_id: str) -> Dict[str, Any]:
    """
    用户对比分析
    
    比较两个用户的行为模式和标签,
    找出相似性和差异性。
    """
    profile1 = analyzer.analyze_behavior_pattern(user1_id)
    profile2 = analyzer.analyze_behavior_pattern(user2_id)
    
    tags1 = await generate_user_tags(user1_id)
    tags2 = await generate_user_tags(user2_id)
    
    # 计算相似度
    common_tags = set(tags1.get("tags", [])) & set(tags2.get("tags", []))
    similarity_score = len(common_tags) / max(len(set(tags1.get("tags", []))), 1)
    
    return {
        "users": [user1_id, user2_id],
        "comparison": {
            "common_tags": list(common_tags),
            "similarity_score": round(similarity_score, 2),
            "differences": {
                "unique_to_user1": list(set(tags1.get("tags", [])) - set(tags2.get("tags", []))),
                "unique_to_user2": list(set(tags2.get("tags", [])) - set(tags1.get("tags", [])))
            }
        },
        "profiles": {
            user1_id: {
                "metrics": profile1.get("metrics", {}),
                "engagement": profile1.get("patterns", {}).get("engagement_score", 0)
            },
            user2_id: {
                "metrics": profile2.get("metrics", {}),
                "engagement": profile2.get("patterns", {}).get("engagement_score", 0)
            }
        }
    }

if __name__ == "__main__":
    # 启动服务
    app.run(transport="sse", host="127.0.0.1", port=8080)

这个用户画像系统展示了Resource Template的强大能力:

  1. 复杂分析逻辑:封装了用户行为分析算法
  2. 多维度数据:支持行为分析、标签生成、用户对比
  3. 参数化查询:支持时间范围、用户对比等参数
  4. 结构化输出:返回JSON格式的详细分析结果

4.3 客户端集成示例

# profile_client.py
import asyncio
from mcp import ClientSession
from mcp.client.sse import sse_client
from contextlib import AsyncExitStack

class ProfileAnalysisClient:
    """用户画像分析客户端"""
    
    def __init__(self):
        self.exit_stack = AsyncExitStack()
        self.session = None
    
    async def connect(self):
        """连接到画像分析服务"""
        read_stream, write_stream = await self.exit_stack.enter_async_context(
            sse_client("http://127.0.0.1:8080/sse")
        )
        
        self.session = await self.exit_stack.enter_async_context(
            ClientSession(read_stream, write_stream)
        )
        
        await self.session.initialize()
        print("✅ 已连接到用户画像分析服务")
    
    async def analyze_user_behavior(self, user_id: str, days: int = 7):
        """分析用户行为"""
        uri = f"profile://{user_id}/behavior?period={days}"
        response = await self.session.read_resource(uri)
        
        if response.contents:
            result = json.loads(response.contents[0].text)
            print(f"\n📊 用户 {user_id} 的行为分析报告:")
            print(f"   分析周期: {result['metadata']['analysis_period']}")
            print(f"   总会话数: {result['metrics']['total_sessions']}")
            print(f"   购买次数: {result['metrics']['total_purchases']}")
            print(f"   分享次数: {result['metrics']['total_shares']}")
            print(f"   参与度分数: {result['patterns']['engagement_score']}/10")
            print(f"   活跃时间段: {result['patterns']['active_hours']}")
            print(f"   推荐建议: {', '.join(result['recommendations'])}")
    
    async def get_user_tags(self, user_id: str):
        """获取用户标签"""
        uri = f"profile://{user_id}/tags"
        response = await self.session.read_resource(uri)
        
        if response.contents:
            result = json.loads(response.contents[0].text)
            print(f"\n🏷️ 用户 {user_id} 的标签:")
            print(f"   标签数量: {result['tag_count']}")
            print(f"   行为标签: {', '.join(result['tag_details']['behavior_based'])}")
            print(f"   时间标签: {', '.join(result['tag_details']['time_based'])}")
    
    async def compare_users(self, user1_id: str, user2_id: str):
        """比较两个用户"""
        uri = f"profile://compare?user1={user1_id}&user2={user2_id}"
        response = await self.session.read_resource(uri)
        
        if response.contents:
            result = json.loads(response.contents[0].text)
            print(f"\n🔍 用户对比分析:")
            print(f"   相似度分数: {result['comparison']['similarity_score']}")
            print(f"   共同标签: {', '.join(result['comparison']['common_tags'])}")
            print(f"   {user1_id} 独有标签: {', '.join(result['comparison']['differences']['unique_to_user1'])}")
            print(f"   {user2_id} 独有标签: {', '.join(result['comparison']['differences']['unique_to_user2'])}")
    
    async def close(self):
        """关闭连接"""
        await self.exit_stack.aclose()

# 使用示例
async def main():
    client = ProfileAnalysisClient()
    
    try:
        await client.connect()
        
        # 分析用户行为
        await client.analyze_user_behavior("1001", 30)
        await client.analyze_user_behavior("1002", 30)
        
        # 获取用户标签
        await client.get_user_tags("1001")
        await client.get_user_tags("1002")
        
        # 比较用户
        await client.compare_users("1001", "1002")
        
    finally:
        await client.close()

if __name__ == "__main__":
    asyncio.run(main())

这个客户端展示了如何与复杂的Resource Template服务交互,获取有价值的用户洞察。通过组合不同的Resource Template,可以构建出强大的用户分析功能。

5. 性能优化与最佳实践

在真实生产环境中使用Resource Template时,性能优化至关重要。以下是我在实践中总结的一些最佳实践。

5.1 缓存策略优化

Resource Template服务可能会被频繁调用,合理的缓存策略可以显著提升性能:

import time
from functools import lru_cache
from typing import Dict, Any
import hashlib

class ResourceCache:
    """资源缓存管理器"""
    
    def __init__(self, max_size: int = 1000, ttl: int = 300):
        """
        初始化缓存
        
        Args:
            max_size: 最大缓存条目数
            ttl: 缓存存活时间(秒)
        """
        self.max_size = max_size
        self.ttl = ttl
        self._cache: Dict[str, Dict[str, Any]] = {}
        self._timestamps: Dict[str, float] = {}
    
    def _generate_key(self, uri_template: str, **kwargs) -> str:
        """生成缓存键"""
        # 基于URI模板和参数生成唯一键
        key_data = f"{uri_template}:{sorted(kwargs.items())}"
        return hashlib.md5(key_data.encode()).hexdigest()
    
    def get(self, uri_template: str, **kwargs) -> Any:
        """获取缓存值"""
        key = self._generate_key(uri_template, **kwargs)
        
        if key in self._cache:
            # 检查是否过期
            if time.time() - self._timestamps[key] < self.ttl:
                return self._cache[key]
            else:
                # 清理过期缓存
                del self._cache[key]
                del self._timestamps[key]
        
        return None
    
    def set(self, uri_template: str, value: Any, **kwargs):
        """设置缓存值"""
        key = self._generate_key(uri_template, **kwargs)
        
        # 清理最旧的条目如果达到最大大小
        if len(self._cache) >= self.max_size:
            oldest_key = min(self._timestamps.items(), key=lambda x: x[1])[0]
            del self._cache[oldest_key]
            del self._timestamps[oldest_key]
        
        self._cache[key] = value
        self._timestamps[key] = time.time()
    
    def clear(self):
        """清空缓存"""
        self._cache.clear()
        self._timestamps.clear()

# 在Resource Template中使用缓存
cache = ResourceCache(max_size=500, ttl=60)  # 缓存500条,60秒过期

@app.resource(
    uri="data://{dataset}/query?filters={filters}",
    name="cached_data_query",
    description="带缓存的数据查询\n:param dataset: 数据集名称\n:param filters: 过滤条件",
    mime_type="application/json"
)
async def query_cached_data(dataset: str, filters: str) -> Dict[str, Any]:
    """
    带缓存的数据查询
    
    对频繁查询的数据进行缓存,
    减少数据库或API调用压力。
    """
    # 检查缓存
    cached_result = cache.get("data://{dataset}/query?filters={filters}", 
                             dataset=dataset, filters=filters)
    
    if cached_result:
        cached_result["cached"] = True
        cached_result["cached_at"] = time.time()
        return cached_result
    
    # 模拟数据查询(实际项目中这里可能是数据库查询或API调用)
    await asyncio.sleep(0.5)  # 模拟查询延迟
    
    result = {
        "dataset": dataset,
        "filters": filters,
        "data": [{"id": i, "value": f"data_{i}"} for i in range(10)],
        "timestamp": time.time(),
        "cached": False
    }
    
    # 设置缓存
    cache.set("data://{dataset}/query?filters={filters}", result, 
              dataset=dataset, filters=filters)
    
    return result

5.2 异步处理优化

对于耗时的Resource Template处理,使用异步可以显著提升并发性能:

import asyncio
from concurrent.futures import ThreadPoolExecutor
from mcp.server.fastmcp import FastMCP

app = FastMCP("async-resource-service")

# 创建线程池用于CPU密集型操作
executor = ThreadPoolExecutor(max_workers=4)

@app.resource(
    uri="process://heavy/{task_id}",
    name="heavy_processing",
    description="执行重型计算任务\n:param task_id: 任务ID",
    mime_type="application/json"
)
async def heavy_processing(task_id: str) -> Dict[str, Any]:
    """
    重型计算任务
    
    使用线程池避免阻塞事件循环,
    适合CPU密集型操作。
    """
    def cpu_intensive_task():
        """模拟CPU密集型任务"""
        import time
        result = 0
        for i in range(10000000):
            result += i * i
        return result
    
    # 在线程池中执行CPU密集型任务
    loop = asyncio.get_event_loop()
    result = await loop.run_in_executor(executor, cpu_intensive_task)
    
    return {
        "task_id": task_id,
        "result": result,
        "status": "completed",
        "processed_at": time.time()
    }

@app.resource(
    uri="aggregate://{operation}/{data_ids}",
    name="data_aggregation",
    description="数据聚合操作\n:param operation: 聚合操作类型\n:param data_ids: 数据ID列表(逗号分隔)",
    mime_type="application/json"
)
async def aggregate_data(operation: str, data_ids: str) -> Dict[str, Any]:
    """
    并行数据聚合
    
    使用asyncio.gather并行处理多个数据源,
    适合IO密集型操作。
    """
    ids = [id.strip() for id in data_ids.split(',')]
    
    # 定义异步数据获取函数
    async def fetch_data(data_id: str):
        """模拟异步数据获取"""
        await asyncio.sleep(0.1)  # 模拟网络延迟
        return {
            "id": data_id,
            "value": hash(data_id) % 100  # 模拟数据值
        }
    
    # 并行获取所有数据
    tasks = [fetch_data(data_id) for data_id in ids]
    data_items = await asyncio.gather(*tasks)
    
    # 根据操作类型进行聚合
    values = [item["value"] for item in data_items]
    
    if operation == "sum":
        result = sum(values)
    elif operation == "avg":
        result = sum(values) / len(values) if values else 0
    elif operation == "max":
        result = max(values) if values else 0
    elif operation == "min":
        result = min(values) if values else 0
    else:
        result = None
    
    return {
        "operation": operation,
        "data_ids": ids,
        "result": result,
        "items_processed": len(data_items),
        "details": data_items
    }

5.3 监控与日志

在生产环境中,完善的监控和日志是必不可少的:

import logging
from datetime import datetime
from typing import Dict, Any
from mcp.server.fastmcp import FastMCP

# 配置日志
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler('mcp_server.log'),
        logging.StreamHandler()
    ]
)

logger = logging.getLogger(__name__)

app = FastMCP("monitored-resource-service")

class ResourceMonitor:
    """资源使用监控器"""
    
    def __init__(self):
        self.stats = {
            "total_requests": 0,
            "successful_requests": 0,
            "failed_requests": 0,
            "average_response_time": 0,
            "requests_by_resource": {},
            "errors_by_type": {}
        }
        self.start_time = datetime.now()
    
    def record_request(self, resource_name: str, duration: float, success: bool = True, error: str = None):
        """记录请求统计"""
        self.stats["total_requests"] += 1
        
        if success:
            self.stats["successful_requests"] += 1
        else:
            self.stats["failed_requests"] += 1
            if error:
                self.stats["errors_by_type"][error] = self.stats["errors_by_type"].get(error, 0) + 1
        
        # 更新平均响应时间
        current_avg = self.stats["average_response_time"]
        total_req = self.stats["total_requests"]
        self.stats["average_response_time"] = (current_avg * (total_req - 1) + duration) / total_req
        
        # 按资源统计
        if resource_name not in self.stats["requests_by_resource"]:
            self.stats["requests_by_resource"][resource_name] = {
                "count": 0,
                "total_time": 0,
                "avg_time": 0
            }
        
        resource_stats = self.stats["requests_by_resource"][resource_name]
        resource_stats["count"] += 1
        resource_stats["total_time"] += duration
        resource_stats["avg_time"] = resource_stats["total_time"] / resource_stats["count"]
    
    def get_stats(self) -> Dict[str, Any]:
        """获取统计信息"""
        uptime = (datetime.now() - self.start_time).total_seconds()
        
        return {
            **self.stats,
            "uptime_seconds": uptime,
            "requests_per_second": self.stats["total_requests"] / uptime if uptime > 0 else 0,
            "success_rate": (self.stats["successful_requests"] / self.stats["total_requests"] * 100 
                           if self.stats["total_requests"] > 0 else 0)
        }

monitor = ResourceMonitor()

# 装饰器:监控Resource Template性能
def monitor_resource(func):
    """资源监控装饰器"""
    async def wrapper(*args, **kwargs):
        start_time = datetime.now()
        resource_name = func.__name__
        
        try:
            result = await func(*args, **kwargs)
            duration = (datetime.now() - start_time).total_seconds()
            
            # 记录成功请求
            monitor.record_request(resource_name, duration, success=True)
            logger.info(f"Resource '{resource_name}' executed in {duration:.3f}s")
            
            return result
            
        except Exception as e:
            duration = (datetime.now() - start_time).total_seconds()
            
            # 记录失败请求
            monitor.record_request(resource_name, duration, success=False, error=str(e))
            logger.error(f"Resource '{resource_name}' failed after {duration:.3f}s: {str(e)}")
            
            # 返回错误信息
            return {
                "error": str(e),
                "resource": resource_name,
                "timestamp": datetime.now().isoformat()
            }
    
    return wrapper

@app.resource(
    uri="monitored://{resource_name}",
    name="monitored_resource",
    description="监控示例资源\n:param resource_name: 资源名称",
    mime_type="application/json"
)
@monitor_resource
async def monitored_resource_example(resource_name: str) -> Dict[str, Any]:
    """
    监控示例资源
    
    展示如何监控Resource Template的性能和错误情况。
    """
    # 模拟一些处理
    await asyncio.sleep(0.1)
    
    if resource_name == "error":
        raise ValueError("模拟错误情况")
    
    return {
        "resource_name": resource_name,
        "status": "success",
        "processed_at": datetime.now().isoformat(),
        "data": {"sample": "data", "value": 42}
    }

@app.resource(
    uri="stats://system",
    name="system_stats",
    description="获取系统统计信息",
    mime_type="application/json"
)
async def get_system_stats() -> Dict[str, Any]:
    """
    系统统计信息
    
    返回Resource Template服务的运行统计,
    用于监控和调试。
    """
    stats = monitor.get_stats()
    
    # 添加系统信息
    import psutil
    import os
    
    process = psutil.Process(os.getpid())
    
    return {
        "monitoring": stats,
        "system": {
            "cpu_percent": process.cpu_percent(),
            "memory_mb": process.memory_info().rss / 1024 / 1024,
            "thread_count": process.num_threads(),
            "connections": len(process.connections()) if hasattr(process, 'connections') else 0
        },
        "timestamp": datetime.now().isoformat()
    }

5.4 安全最佳实践

安全是生产环境中的重要考虑因素:

import re
from typing import Optional
from mcp.server.fastmcp import FastMCP

app = FastMCP("secure-resource-service")

class SecurityValidator:
    """安全验证器"""
    
    @staticmethod
    def validate_user_id(user_id: str) -> bool:
        """验证用户ID格式"""
        # 只允许字母、数字、下划线和短横线
        pattern = r'^[a-zA-Z0-9_-]{1,50}$'
        return bool(re.match(pattern, user_id))
    
    @staticmethod
    def sanitize_input(input_str: str) -> str:
        """清理输入字符串"""
        # 移除潜在的恶意字符
        sanitized = re.sub(r'[<>"\']', '', input_str)
        # 限制长度
        return sanitized[:1000]
    
    @staticmethod
    def validate_query_params(params: Dict[str, str]) -> Optional[str]:
        """验证查询参数"""
        for key, value in params.items():
            # 检查参数名是否合法
            if not re.match(r'^[a-zA-Z_][a-zA-Z0-9_]*$', key):
                return f"Invalid parameter name: {key}"
            
            # 检查参数值长度
            if len(value) > 1000:
                return f"Parameter value too long: {key}"
        
        return None

validator = SecurityValidator()

@app.resource(
    uri="secure://user/{user_id}/data",
    name="secure_user_data",
    description="安全用户数据访问\n:param user_id: 用户ID",
    mime_type="application/json"
)
async def get_secure_user_data(user_id: str) -> Dict[str, Any]:
    """
    安全用户数据访问
    
    包含输入验证和清理的安全示例。
    """
    # 验证用户ID
    if not validator.validate_user_id(user_id):
        return {
            "error": "Invalid user ID format",
            "user_id": user_id,
            "allowed_pattern": "字母、数字、下划线、短横线,1-50字符"
        }
    
    # 清理用户ID
    sanitized_id = validator.sanitize_input(user_id)
    
    # 模拟数据访问(实际项目中这里可能是数据库查询)
    user_data = {
        "user_id": sanitized_id,
        "username": f"user_{sanitized_id}",
        "access_level": "standard",
        "last_accessed": datetime.now().isoformat()
    }
    
    # 记录访问日志(不含敏感信息)
    logger.info(f"Secure data access for user: {sanitized_id}")
    
    return user_data

@app.resource(
    uri="secure://query?{query_string}",
    name="secure_query",
    description="安全查询接口\n支持多个查询参数",
    mime_type="application/json"
)
async def secure_query_endpoint(**kwargs) -> Dict[str, Any]:
    """
    安全查询接口
    
    演示如何处理动态查询参数的安全验证。
    """
    # 验证查询参数
    validation_error = validator.validate_query_params(kwargs)
    if validation_error:
        return {"error": validation_error, "params": kwargs}
    
    # 清理所有参数值
    sanitized_params = {
        key: validator.sanitize_input(value)
        for key, value in kwargs.items()
    }
    
    # 模拟查询处理
    result = {
        "query_params": sanitized_params,
        "result_count": len(sanitized_params),
        "processed_at": datetime.now().isoformat(),
        "security_check": "passed"
    }
    
    # 记录查询(不含参数值)
    logger.info(f"Secure query with {len(sanitized_params)} parameters")
    
    return result

# 速率限制装饰器
from functools import wraps
import time

class RateLimiter:
    """简单的速率限制器"""
    
    def __init__(self, max_calls: int, period: float):
        self.max_calls = max_calls
        self.period = period
        self.calls = []
    
    def __call__(self, func):
        @wraps(func)
        async def wrapper(*args, **kwargs):
            now = time.time()
            
            # 清理过期的调用记录
            self.calls = [call_time for call_time in self.calls 
                         if now - call_time < self.period]
            
            # 检查是否超过限制
            if len(self.calls) >= self.max_calls:
                wait_time = self.period - (now - self.calls[0])
                return {
                    "error": f"Rate limit exceeded. Try again in {wait_time:.1f} seconds.",
                    "retry_after": wait_time
                }
            
            # 记录本次调用
            self.calls.append(now)
            
            # 执行原函数
            return await func(*args, **kwargs)
        
        return wrapper

# 创建速率限制器:每分钟最多10次调用
rate_limiter = RateLimiter(max_calls=10, period=60)

@app.resource(
    uri="limited://resource",
    name="rate_limited_resource",
    description="速率限制的资源示例",
    mime_type="application/json"
)
@rate_limiter
async def rate_limited_resource() -> Dict[str, Any]:
    """
    速率限制的资源
    
    演示如何实现基本的速率限制。
    """
    return {
        "message": "This is a rate-limited resource",
        "remaining_calls": rate_limiter.max_calls - len(rate_limiter.calls),
        "reset_in": rate_limiter.period - (time.time() - rate_limiter.calls[0]) 
                   if rate_limiter.calls else 0
    }

这些最佳实践涵盖了缓存、异步处理、监控、安全等多个方面,能够帮助你在生产环境中构建稳定、高效、安全的Resource Template服务。

在实际项目中,我发现最关键的几点是:合理的缓存策略可以大幅提升性能,完善的监控可以帮助快速定位问题,而严格的安全验证则是防止攻击的第一道防线。特别是参数验证,很多安全问题都源于对用户输入过于信任。

Resource Template的真正价值在于它的灵活性——你可以根据实际需求组合不同的优化策略。比如对于查询频繁但变化不频繁的数据,使用较长的TTL缓存;对于安全性要求高的接口,实施严格的输入验证和速率限制;对于计算密集型的操作,使用异步处理避免阻塞。

记住,没有一种方案适合所有场景。最好的做法是根据具体的业务需求和数据特性,选择合适的优化策略。通过监控系统观察实际运行情况,不断调整和优化,才能构建出既高效又可靠的Resource Template服务。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐