⚡ 模块2:高阶能力实战

2.1 向量数据库对接(Chroma)

在第三阶段项目1中,我们用 NumPy 手写了简易向量存储。生产环境需要专业的向量数据库。Chroma 是最适合入门者的向量数据库——Python 原生、无需单独部署、API 简洁。

# 安装 Chroma
# pip install chromadb sentence-transformers

import chromadb
from chromadb.utils import embedding_functions
from pathlib import Path

class ChromaMemoryStore:
    """基于 Chroma 的长期记忆向量存储"""

    def __init__(self, persist_dir: str = "./agent_chroma_db"):
        self.client = chromadb.PersistentClient(path=persist_dir)

        # 使用 sentence-transformers 作为嵌入函数
        self.embed_fn = embedding_functions.SentenceTransformerEmbeddingFunction(
            model_name="paraphrase-multilingual-MiniLM-L12-v2"
        )

        # 创建/获取集合
        self.collection = self.client.get_or_create_collection(
            name="agent_memory",
            embedding_function=self.embed_fn,
            metadata={"description": "Agent 长期记忆存储"}
        )

        print(f"✅ Chroma 初始化完成,现有 {self.collection.count()} 条记忆")

    def remember(self, content: str, metadata: dict = None) -> str:
        """存储一条记忆"""
        import uuid
        mem_id = str(uuid.uuid4())[:8]

        self.collection.add(
            documents=[content],
            metadatas=[metadata or {}],
            ids=[mem_id]
        )
        return mem_id

    def recall(self, query: str, top_k: int = 5) -> list[dict]:
        """语义搜索相关记忆"""
        results = self.collection.query(
            query_texts=[query],
            n_results=top_k,
            include=["documents", "metadatas", "distances"]
        )

        memories = []
        if results["documents"] and results["documents"][0]:
            for i, doc in enumerate(results["documents"][0]):
                memories.append({
                    "content": doc,
                    "metadata": results["metadatas"][0][i] if results["metadatas"] else {},
                    "relevance": 1 - results["distances"][0][i],  # 距离→相似度
                })
        return memories

    def forget(self, memory_id: str):
        """删除指定记忆"""
        self.collection.delete(ids=[memory_id])

    def clear(self):
        """清空所有记忆"""
        self.collection.delete(where={})  # 需要先有数据才能清空
        # 删除并重建
        self.client.delete_collection("agent_memory")
        self.collection = self.client.create_collection(
            name="agent_memory",
            embedding_function=self.embed_fn
        )

# 使用示例
if __name__ == "__main__":
    store = ChromaMemoryStore()

    # 存储记忆
    store.remember("用户张三喜欢简洁的回答风格", {"type": "preference", "user": "zhangsan"})
    store.remember("项目A的截止日期是6月30日", {"type": "fact", "project": "A"})
    store.remember("上次讨论了Python性能优化方案", {"type": "history", "topic": "python"})

    # 语义搜索
    results = store.recall("用户的回答偏好是什么?")
    for r in results:
        print(f"  [{r['relevance']:.3f}] {r['content']}")
    # 能准确找到"喜欢简洁风格",即使没有关键词完全匹配

2.2 复杂多任务并行执行

#!/usr/bin/env python3
"""ParallelAgent — 多任务并行执行 Agent"""

import asyncio
from concurrent.futures import ThreadPoolExecutor, as_completed
from anthropic import Anthropic
import os, json, time
from dotenv import load_dotenv
load_dotenv()

class ParallelTaskExecutor:
    """并行任务执行器"""

    def __init__(self, max_workers: int = 4):
        self.client = Anthropic(api_key=os.getenv("ANTHROPIC_API_KEY"))
        self.model = "claude-sonnet-4-6"
        self.max_workers = max_workers

    def execute_parallel(self, tasks: list[dict]) -> list[dict]:
        """
        并行执行多个独立任务
        tasks: [{"id": "1", "description": "任务描述", "context": "上下文"}, ...]
        """
        start_time = time.time()
        results = []

        with ThreadPoolExecutor(max_workers=self.max_workers) as executor:
            # 提交所有任务
            future_to_task = {
                executor.submit(self._execute_single, task): task
                for task in tasks
            }

            # 收集结果
            for future in as_completed(future_to_task):
                task = future_to_task[future]
                try:
                    result = future.result(timeout=60)
                    results.append({"id": task["id"], "status": "done", "result": result})
                except Exception as e:
                    results.append({"id": task["id"], "status": "failed", "error": str(e)})

        elapsed = time.time() - start_time
        print(f"⚡ 并行执行 {len(tasks)} 个任务,总耗时 {elapsed:.1f}s")

        # 按ID排序
        results.sort(key=lambda x: x["id"])
        return results

    def _execute_single(self, task: dict) -> str:
        """执行单个任务"""
        prompt = f"请完成以下任务(直接输出结果,不要解释过程):\n\n{task.get('description', '')}"
        if task.get("context"):
            prompt += f"\n\n上下文: {task['context']}"

        response = self.client.messages.create(
            model=self.model, max_tokens=500,
            messages=[{"role": "user", "content": prompt}],
            temperature=0.3)
        return response.content[0].text.strip()

    def execute_with_priority(self, tasks: list[dict]) -> list[dict]:
        """
        带优先级的任务调度
        tasks 中每个包含 "priority": 1-5 (1最高)
        """
        # 按优先级排序
        sorted_tasks = sorted(tasks, key=lambda t: t.get("priority", 3))
        high_priority = [t for t in sorted_tasks if t.get("priority", 3) <= 2]
        low_priority = [t for t in sorted_tasks if t.get("priority", 3) > 2]

        print(f"📊 高优先级: {len(high_priority)} 个, 低优先级: {len(low_priority)} 个")

        # 先执行高优先级,再执行低优先级
        results = []
        if high_priority:
            results.extend(self.execute_parallel(high_priority))
        if low_priority:
            results.extend(self.execute_parallel(low_priority))
        return results


# 测试
if __name__ == "__main__":
    executor = ParallelTaskExecutor()

    tasks = [
        {"id": "1", "description": "用一句话总结Python的优点", "priority": 1},
        {"id": "2", "description": "列出3个常见的Git命令", "priority": 3},
        {"id": "3", "description": "解释什么是RESTful API(50字以内)", "priority": 2},
        {"id": "4", "description": "写一个Python列表推导式示例", "priority": 3},
    ]

    print("⚡ 并行执行模式:")
    results = executor.execute_parallel(tasks)
    for r in results:
        print(f"  任务{r['id']}: {r.get('result', r.get('error', ''))[:80]}...")

    print("\n📊 优先级调度模式:")
    results = executor.execute_with_priority(tasks)
    for r in results:
        print(f"  任务{r['id']}: {r.get('result', r.get('error', ''))[:80]}...")

2.3 幻觉抑制与输出一致性

策略实现方法效果
引用约束 Prompt: "只使用提供的信息。不确定时明确说'不确定'" ⭐⭐⭐⭐ 大幅减少编造
多路验证 同一问题问3次(temperature不同),取交集 ⭐⭐⭐ 过滤偶发幻觉
事实核查Agent 独立Agent专门验证输出中的事实陈述 ⭐⭐⭐⭐⭐ 工业级方案
输出格式约束 要求结构化输出 + 每项标注确定性(确定/可能/不确定) ⭐⭐⭐⭐ 用户可判断可信度
知识边界声明 Agent 在回答前声明信息来源和知识截止日期 ⭐⭐⭐ 设定正确预期
# 幻觉抑制核心 Prompt 模板
HALLUCINATION_CONTROL_PROMPT = """
## 🛡️ 信息准确性守则

1. **区分事实与推测**
   - 确定的事实: 用肯定语气
   - 推测/不确定: 使用"可能"、"据我了解"、"建议核实"
   - 不知道: 明确说"我目前没有这方面的确切信息"

2. **标注信息来源**
   每个关键陈述后标注来源类型:
   [已知事实] [文档依据] [推理结果] [外部知识/请核实]

3. **自我检查清单**
   回答完成后,在脑海中默查:
   ✓ 所有数字是否有依据?
   ✓ 所有"事实"是我知道的还是猜测的?
   ✓ 是否有需要用户进一步核实的建议?

4. **不确定时的标准回复**
   "关于[X],我目前没有足够的信息给出确切答案。
    建议: [提供获取信息的方法/途径]"
"""

def verify_facts(client, text: str) -> dict:
    """用独立Claude调用验证文本中的事实"""
    prompt = f"""分析以下文本中的事实陈述,逐条标注可信度。

文本:
{text}

对每个事实陈述,标注:
- fact: 事实陈述内容
- confidence: high/medium/low
- reason: 判断理由
- needs_verification: true/false

输出JSON数组。"""

    resp = client.messages.create(
        model="claude-sonnet-4-6", max_tokens=500,
        messages=[{"role":"user","content":prompt}], temperature=0.1)
    try:
        text = resp.content[0].text.strip()
        if "```" in text: text = text.split("```")[1].split("```")[0]
        return json.loads(text)
    except:
        return []

✏️ 课后练习

  1. Chroma 实战:将第三阶段 QAAgent 的向量存储替换为 Chroma,对比检索速度和准确率。
  2. 并行性能测试:分别用串行和并行方式执行5个任务,对比总耗时。记录加速比。
  3. 幻觉检测实验:故意让 Agent 回答一个它不可能知道的事实(如"2040年的AI市场规模"),测试幻觉抑制策略的效果。
  4. 优先级调度:完善优先级调度器,增加"抢占"功能——高优先级任务到达时暂停低优先级任务。
← 高阶架构 下一模块:性能优化 →