知识摄取 Pipeline 设计

文档版本: v1.0 日期: 2026-06-09 作者: 数据工程与 AI 系统架构组 定位: 工程实践导向,面向构建个人/团队 AI 知识库的工程师


概述

知识摄取(Knowledge Ingestion)是 AI 知识库系统的”消化系统”——将来自各种信源的原始信息转化为可检索、可推理的结构化知识条目。一个设计良好的摄取 Pipeline 应当具备:

  1. 广覆盖:能从手动笔记、网页、工具输出、AI 会话等多种来源摄入
  2. 高质量:通过 LLM 提炼,将信息噪声过滤,保留核心洞察
  3. 可持续:具备幂等性和自动化触发,无需人工持续介入
  4. 可观测:完整的摄取日志、质量评分、错误追踪

本文给出从架构设计到完整代码实现的端到端工程方案。


1. 知识来源的分类

知识来源按自动化程度分为四个层级,各有其适用场景和处理策略。

1.1 主动输入(Manual)

用户主动编写的内容,质量最高,但产出量最低。

  • 手动写笔记:Obsidian 中直接编辑的 Markdown 文件
  • 分析报告:深度调研文档、技术决策记录(ADR)
  • 代码注释导出:从代码库批量提取 // NOTE: // HACK: 等标注

处理策略:保留原文,仅补充元数据(自动打标签、生成摘要)。

1.2 半自动输入(Semi-auto)

用户有意识收集,但格式散乱,需要处理。

  • 网页收藏:浏览器书签导出、Raindrop/Instapaper 等剪藏工具
  • PDF/文档解析:技术规格书、论文、官方文档
  • 微信/Slack 消息导出:团队讨论中的关键决策片段

处理策略:内容清洗 + LLM 摘要 + 自动分类

1.3 全自动输入(Fully-auto)

系统在后台持续产出,无需人工介入。

  • AI 分析会话提取:每次 Claude/ChatGPT 会话结束后提取经验(本文重点)
  • 日志分析结果:自动化日志扫描发现的异常模式
  • 监控告警归档:Grafana/PagerDuty 告警处理记录

处理策略:触发式摄取 + 结构化提炼 + 置信度评估

1.4 工具输出(Tool Output)

工程工具产生的结构化或半结构化报告。

  • 性能分析报告:Perfetto trace 分析、Simpleperf 火焰图结论
  • 代码审查结论:Code Review 评论、静态分析报告(lint/sonar)
  • 测试报告:崩溃分析、ANR 堆栈摘要

处理策略:格式解析 + 知识点提取 + 关联代码上下文


2. 信息摄取 Agent 设计(Hermes 角色)

2.1 角色定义

在知识库系统的 Agent 体系中,Hermes(信使神)负责信息的搬运与转化。其核心职责如下:

信息源 ──→ Hermes ──→ 知识库
           ┌─────────────────────┐
           │ 1. 监听多个信息源    │
           │ 2. 去重与优先级排序  │
           │ 3. LLM 提炼         │
           │ 4. 写入 Obsidian    │
           │ 5. 更新向量数据库   │
           └─────────────────────┘

2.2 设计原则

幂等性:同一信息摄取多次,结果一致,不产生重复条目。实现方式:对内容计算 SHA-256 指纹,写入前检查是否已存在。

可观测性:每次摄取记录详细日志,包括来源、处理时间、LLM token 消耗、写入结果。

可配置的摄取规则:通过 YAML 配置文件定义每种来源的处理策略,无需修改代码。

2.3 Hermes Agent 核心架构

# hermes/core.py
import hashlib
import json
import logging
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from datetime import datetime
from enum import Enum
from pathlib import Path
from typing import Any, Optional
 
logger = logging.getLogger("hermes")
 
 
class SourceType(Enum):
    MANUAL = "manual"
    WEB = "web"
    AI_SESSION = "ai_session"
    TOOL_OUTPUT = "tool_output"
    DOCUMENT = "document"
 
 
@dataclass
class RawContent:
    """原始信息单元,摄取前的容器"""
    source_type: SourceType
    content: str                          # 原始文本
    source_url: Optional[str] = None      # 来源 URL(如适用)
    source_file: Optional[str] = None     # 来源文件路径
    metadata: dict = field(default_factory=dict)
    ingested_at: datetime = field(default_factory=datetime.utcnow)
 
    @property
    def fingerprint(self) -> str:
        """内容指纹,用于幂等性去重"""
        return hashlib.sha256(self.content.encode("utf-8")).hexdigest()[:16]
 
 
@dataclass
class KnowledgeEntry:
    """标准知识条目,摄取后的输出"""
    title: str
    summary: str
    problem_pattern: str
    solution: str
    evidence: str
    confidence: float            # 0.0 - 1.0
    tags: list[str]
    related_tools: list[str]
    applicability: str
    source_fingerprint: str
    source_type: str
    created_at: datetime = field(default_factory=datetime.utcnow)
    raw_content: Optional[str] = None  # 保留原文
 
    def to_obsidian_markdown(self) -> str:
        """生成 Obsidian 格式 Markdown"""
        tags_str = "\n".join(f"  - {t}" for t in self.tags)
        tools_str = ", ".join(self.related_tools) if self.related_tools else "无"
        return f"""---
title: "{self.title}"
date: {self.created_at.strftime('%Y-%m-%d')}
confidence: {self.confidence:.2f}
source_type: {self.source_type}
fingerprint: {self.source_fingerprint}
tags:
{tags_str}
---
 
# {self.title}
 
## 摘要
{self.summary}
 
## 问题模式
{self.problem_pattern}
 
## 解决方案
{self.solution}
 
## 依据与来源
{self.evidence}
 
## 相关工具
{tools_str}
 
## 适用范围
{self.applicability}
 
---
*置信度: {self.confidence:.0%} | 来源: {self.source_type} | 摄取时间: {self.created_at.strftime('%Y-%m-%d %H:%M')}*
"""
 
 
class BaseIngester(ABC):
    """摄取器基类"""
 
    @abstractmethod
    def can_handle(self, source: RawContent) -> bool:
        pass
 
    @abstractmethod
    async def ingest(self, source: RawContent) -> Optional[KnowledgeEntry]:
        pass

3. 网页内容摄取

3.1 核心流程

URL → Jina Reader → Markdown → 质量过滤 → LLM 摘要 → 知识条目
         ↑
    (动态页面走 Playwright fallback)

3.2 完整实现

# hermes/ingesters/web_ingester.py
import asyncio
import re
import httpx
from typing import Optional
from ..core import BaseIngester, RawContent, KnowledgeEntry, SourceType
 
 
class WebIngester(BaseIngester):
    """网页内容摄取器:Jina Reader + Playwright fallback + LLM 提炼"""
 
    JINA_API = "https://r.jina.ai/"
    MIN_CONTENT_LENGTH = 300   # 过短内容视为低质量
    MAX_CONTENT_LENGTH = 50000  # 避免 token 超限
 
    def __init__(self, llm_client, config: dict = None):
        self.llm = llm_client
        self.config = config or {}
 
    def can_handle(self, source: RawContent) -> bool:
        return source.source_type == SourceType.WEB and source.source_url is not None
 
    async def fetch_via_jina(self, url: str) -> Optional[str]:
        """通过 Jina Reader 将网页转为 Markdown"""
        jina_url = f"{self.JINA_API}{url}"
        headers = {
            "Accept": "text/markdown",
            "X-Return-Format": "markdown",
        }
        # 如果有 Jina API Key,设置更高速率限制
        if api_key := self.config.get("jina_api_key"):
            headers["Authorization"] = f"Bearer {api_key}"
 
        try:
            async with httpx.AsyncClient(timeout=30.0) as client:
                resp = await client.get(jina_url, headers=headers)
                if resp.status_code == 200:
                    return resp.text
        except Exception as e:
            print(f"[WebIngester] Jina fetch failed: {e}")
        return None
 
    async def fetch_via_playwright(self, url: str) -> Optional[str]:
        """Playwright fallback,处理动态渲染页面"""
        try:
            from playwright.async_api import async_playwright
            import html2text
 
            async with async_playwright() as p:
                browser = await p.chromium.launch(headless=True)
                page = await browser.new_page()
                await page.goto(url, wait_until="networkidle", timeout=20000)
                # 等待主要内容加载
                await asyncio.sleep(2)
                html = await page.content()
                await browser.close()
 
            converter = html2text.HTML2Text()
            converter.ignore_links = False
            converter.ignore_images = True
            return converter.handle(html)
        except Exception as e:
            print(f"[WebIngester] Playwright fetch failed: {e}")
            return None
 
    def clean_content(self, raw_markdown: str) -> str:
        """内容清洗:去广告、导航栏、页脚等噪声"""
        lines = raw_markdown.split("\n")
        cleaned = []
        noise_patterns = [
            r"^(Subscribe|Sign up|Newsletter|Advertisement|Related articles)",
            r"^\[.*?(Cookie|Privacy|Terms|Copyright).*?\]",
            r"^(Share|Tweet|Like|Follow us)",
            r"^!\[.*?\]\(.*?ad.*?\)",  # 广告图片
        ]
        noise_re = [re.compile(p, re.IGNORECASE) for p in noise_patterns]
 
        for line in lines:
            if any(p.match(line.strip()) for p in noise_re):
                continue
            cleaned.append(line)
 
        content = "\n".join(cleaned)
        # 合并多余空行
        content = re.sub(r"\n{3,}", "\n\n", content)
        # 截断超长内容
        if len(content) > self.MAX_CONTENT_LENGTH:
            content = content[:self.MAX_CONTENT_LENGTH] + "\n\n[内容已截断...]"
        return content.strip()
 
    def is_quality_content(self, content: str) -> bool:
        """内容质量过滤"""
        if len(content) < self.MIN_CONTENT_LENGTH:
            return False
        # 中英文有效字符占比需超过 50%
        valid_chars = len(re.findall(r'[\w一-鿿]', content))
        total_chars = max(len(content), 1)
        return valid_chars / total_chars > 0.3
 
    async def extract_knowledge(self, url: str, content: str) -> KnowledgeEntry:
        """调用 LLM 从网页内容提炼知识条目"""
        prompt = f"""你是一个知识提炼专家。从以下网页内容中提炼核心知识,输出严格的 JSON 格式。
 
来源 URL: {url}
 
网页内容(Markdown):
{content[:8000]}
 
请提炼并输出以下 JSON(不要包含任何额外文字):
{{
  "title": "简洁标题(15字以内)",
  "summary": "一句话概括核心价值(50字以内)",
  "problem_pattern": "这个知识解决什么问题/适用什么场景",
  "solution": "核心方法/方案/结论(具体可操作)",
  "evidence": "关键数据/引用/来源背书",
  "confidence": 0.75,
  "tags": ["标签1", "标签2", "标签3"],
  "related_tools": ["工具1", "工具2"],
  "applicability": "适用范围和限制条件"
}}"""
 
        response = await self.llm.chat(prompt)
        data = self._parse_llm_json(response)
        return KnowledgeEntry(
            source_fingerprint="",  # 由调用方填充
            source_type=SourceType.WEB.value,
            raw_content=content[:2000],
            **data
        )
 
    def _parse_llm_json(self, response: str) -> dict:
        """解析 LLM 返回的 JSON,带容错处理"""
        import json
        # 提取 JSON 块
        match = re.search(r'\{.*\}', response, re.DOTALL)
        if match:
            try:
                return json.loads(match.group())
            except json.JSONDecodeError:
                pass
        # fallback:返回基础结构
        return {
            "title": "网页内容摘要",
            "summary": response[:100],
            "problem_pattern": "待完善",
            "solution": response[:500],
            "evidence": "网页来源",
            "confidence": 0.5,
            "tags": ["待分类"],
            "related_tools": [],
            "applicability": "待完善"
        }
 
    async def ingest(self, source: RawContent) -> Optional[KnowledgeEntry]:
        url = source.source_url
 
        # Step 1: 获取内容(Jina 优先,Playwright fallback)
        content = await self.fetch_via_jina(url)
        if not content or not self.is_quality_content(content):
            print(f"[WebIngester] Jina 内容不足,尝试 Playwright: {url}")
            content = await self.fetch_via_playwright(url)
 
        if not content or not self.is_quality_content(content):
            print(f"[WebIngester] 内容质量不足,跳过: {url}")
            return None
 
        # Step 2: 清洗内容
        cleaned = self.clean_content(content)
 
        # Step 3: LLM 提炼
        entry = await self.extract_knowledge(url, cleaned)
        entry.source_fingerprint = source.fingerprint
        return entry

4. 从 AI 分析会话提取经验

这是整个摄取 Pipeline 中价值最高的来源。AI 会话包含了问题背景、探索过程、最终结论,是最接近”有效知识”的原始材料。

4.1 触发时机设计

会话结束
   ↓
判断:问题是否已解决?置信度如何?
   ↓
触发经验提取 Prompt(异步,不阻塞会话)
   ↓
LLM 识别:问题模式 / 解决方案 / 关键洞察
   ↓
写入知识库

4.2 经验提取 Prompt 设计

这是最关键的 Prompt,需要引导 LLM 提取有价值的工程经验:

SESSION_EXTRACT_PROMPT = """你是一个工程经验提炼专家。
 
以下是一段技术分析会话记录。请分析并提炼其中的**可复用工程经验**。
 
重要判断标准:
- 如果会话最终找到了明确的解决方案,confidence 应为 0.8-1.0
- 如果只是探索方向,没有明确结论,confidence 应为 0.4-0.6
- 如果会话中途放弃或结论矛盾,confidence 应为 0.2-0.4
 
会话内容:
{session_content}
 
请输出 JSON(仅输出 JSON,不要有其他内容):
{{
  "title": "描述核心问题的简洁标题(20字以内)",
  "summary": "这次分析最重要的发现是什么(1-2句话)",
  "problem_pattern": "触发这类问题的条件/特征/现象(具体描述,帮助以后识别同类问题)",
  "solution": "解决方案或处理方法(可操作的步骤或命令)",
  "evidence": "分析中找到的关键证据(日志片段/代码行/数据指标)",
  "confidence": 0.85,
  "tags": ["领域标签", "技术标签", "问题类型标签"],
  "related_tools": ["使用的工具"],
  "applicability": "这个结论在什么条件下适用,有什么限制",
  "anti_patterns": "分析过程中发现的错误做法(可选)",
  "follow_up": "还有哪些未解决的疑问(可选)"
}}
"""

4.3 完整实现

# hermes/ingesters/session_ingester.py
import json
import re
from typing import Optional
from ..core import BaseIngester, RawContent, KnowledgeEntry, SourceType
 
 
class AISessionIngester(BaseIngester):
    """从 AI 分析会话提取工程经验"""
 
    # 最短会话长度(避免提取无意义的短对话)
    MIN_SESSION_LENGTH = 500
    # 会话中有"解决"等词汇时提升置信度
    RESOLUTION_KEYWORDS = [
        "问题已解决", "找到原因", "根因是", "确认是", "已验证",
        "root cause", "confirmed", "fixed", "resolved"
    ]
 
    def __init__(self, llm_client):
        self.llm = llm_client
 
    def can_handle(self, source: RawContent) -> bool:
        return source.source_type == SourceType.AI_SESSION
 
    def has_resolution(self, content: str) -> bool:
        """启发式判断会话是否找到解决方案"""
        content_lower = content.lower()
        return any(kw.lower() in content_lower for kw in self.RESOLUTION_KEYWORDS)
 
    def preprocess_session(self, content: str) -> str:
        """预处理会话内容:移除无关信息,保留关键交互"""
        lines = content.split("\n")
        processed = []
        for line in lines:
            # 跳过系统提示、工具调用细节等噪声
            if line.startswith("```tool_call") or line.startswith("```tool_result"):
                continue
            # 保留用户问题和助手回答
            processed.append(line)
        return "\n".join(processed)
 
    async def extract_experience(self, session_content: str) -> dict:
        """调用 LLM 提取经验"""
        from .web_ingester import SESSION_EXTRACT_PROMPT
        prompt = SESSION_EXTRACT_PROMPT.format(
            session_content=session_content[:12000]
        )
        response = await self.llm.chat(prompt)
 
        # 解析 JSON
        match = re.search(r'\{.*\}', response, re.DOTALL)
        if match:
            try:
                return json.loads(match.group())
            except json.JSONDecodeError:
                pass
        return None
 
    def adjust_confidence(self, data: dict, session_content: str) -> dict:
        """基于启发式规则调整置信度"""
        base_confidence = data.get("confidence", 0.7)
 
        # 会话中有明确解决标志,提升置信度
        if self.has_resolution(session_content):
            base_confidence = min(1.0, base_confidence + 0.1)
 
        # follow_up 非空说明仍有疑问,适当降低
        if data.get("follow_up") and len(data["follow_up"]) > 20:
            base_confidence = max(0.1, base_confidence - 0.1)
 
        data["confidence"] = round(base_confidence, 2)
        return data
 
    async def ingest(self, source: RawContent) -> Optional[KnowledgeEntry]:
        if len(source.content) < self.MIN_SESSION_LENGTH:
            print("[SessionIngester] 会话内容过短,跳过")
            return None
 
        processed = self.preprocess_session(source.content)
        data = await self.extract_experience(processed)
 
        if not data:
            print("[SessionIngester] LLM 提取失败")
            return None
 
        data = self.adjust_confidence(data, source.content)
 
        # 提取可选字段
        extra_meta = {}
        for key in ("anti_patterns", "follow_up"):
            if data.get(key):
                extra_meta[key] = data.pop(key)
 
        entry = KnowledgeEntry(
            source_fingerprint=source.fingerprint,
            source_type=SourceType.AI_SESSION.value,
            raw_content=source.content[:1000],
            **{k: data[k] for k in [
                "title", "summary", "problem_pattern", "solution",
                "evidence", "confidence", "tags", "related_tools", "applicability"
            ] if k in data}
        )
 
        # 将额外元数据记录到 evidence 字段
        if extra_meta:
            entry.evidence += f"\n\n**补充**:\n{json.dumps(extra_meta, ensure_ascii=False, indent=2)}"
 
        return entry

5. 结构化知识条目生成

5.1 核心 Prompt 设计原则

好的提炼 Prompt 应当:

  1. 明确输出格式:严格 JSON,字段有说明
  2. 引导高密度信息:不要泛泛而谈,要具体可操作
  3. 包含质量控制提示:告知 LLM 如何评估置信度
  4. 容错处理:对 LLM 输出进行健壮性解析

5.2 通用知识提炼 Prompt

KNOWLEDGE_EXTRACTION_PROMPT = """你是一个工程知识提炼专家。将以下原始内容转化为标准知识条目。
 
原始内容类型: {content_type}
原始内容:
{raw_content}
 
评分标准:
- confidence 1.0: 有明确实验数据或权威来源支撑
- confidence 0.7-0.9: 有逻辑推导或多来源印证
- confidence 0.4-0.6: 推测性结论或单一来源
- confidence 0.1-0.3: 不确定或存在争议
 
严格输出以下 JSON(不要包含任何说明文字):
{{
  "title": "15字以内的简洁标题",
  "summary": "一句话概括,让人立刻明白有什么用",
  "problem_pattern": "什么情况下会遇到这个问题(特征、触发条件)",
  "solution": "解决方案(步骤、命令、代码片段——越具体越好)",
  "evidence": "支撑结论的关键依据(数据/日志/源码行)",
  "confidence": 0.0,
  "tags": ["最多5个标签"],
  "related_tools": ["涉及的工具/框架/库"],
  "applicability": "适用范围(版本、平台、条件)和已知限制"
}}
"""

5.3 通用提炼器实现

# hermes/ingesters/generic_ingester.py
import json
import re
from typing import Optional
from ..core import BaseIngester, RawContent, KnowledgeEntry, SourceType
from .prompts import KNOWLEDGE_EXTRACTION_PROMPT
 
 
CONTENT_TYPE_NAMES = {
    SourceType.MANUAL: "手动笔记",
    SourceType.WEB: "网页文章",
    SourceType.AI_SESSION: "AI 分析会话",
    SourceType.TOOL_OUTPUT: "工具输出报告",
    SourceType.DOCUMENT: "技术文档",
}
 
 
class GenericIngester(BaseIngester):
    """通用知识提炼器,处理无专门摄取器的来源"""
 
    def __init__(self, llm_client):
        self.llm = llm_client
 
    def can_handle(self, source: RawContent) -> bool:
        return True  # fallback,处理所有类型
 
    async def ingest(self, source: RawContent) -> Optional[KnowledgeEntry]:
        content_type = CONTENT_TYPE_NAMES.get(source.source_type, "未知类型")
        prompt = KNOWLEDGE_EXTRACTION_PROMPT.format(
            content_type=content_type,
            raw_content=source.content[:10000]
        )
 
        response = await self.llm.chat(prompt)
        match = re.search(r'\{.*\}', response, re.DOTALL)
        if not match:
            return None
 
        try:
            data = json.loads(match.group())
        except json.JSONDecodeError:
            return None
 
        return KnowledgeEntry(
            source_fingerprint=source.fingerprint,
            source_type=source.source_type.value,
            raw_content=source.content[:500],
            **data
        )

6. 去重与合并逻辑

6.1 语义去重策略

简单的文本哈希无法检测语义重复(同一知识用不同表述写了两遍)。需要基于向量相似度进行语义去重。

新条目向量化
     ↓
查询最近邻(TopK=5)
     ↓
相似度 > 0.92 → 视为重复
0.75 < 相似度 < 0.92 → 部分重叠,进行合并
相似度 < 0.75 → 全新知识,直接写入

6.2 完整去重与合并实现

# hermes/deduplicator.py
import json
import logging
from dataclasses import dataclass
from enum import Enum
from typing import Optional
import chromadb
from ..core import KnowledgeEntry
 
logger = logging.getLogger("hermes.dedup")
 
 
class DedupResult(Enum):
    NEW = "new"               # 全新知识,直接写入
    DUPLICATE = "duplicate"   # 完全重复,跳过
    PARTIAL = "partial"       # 部分重叠,需要合并
    CONFLICT = "conflict"     # 矛盾信息,保留双方并标记
 
 
@dataclass
class DedupDecision:
    result: DedupResult
    existing_id: Optional[str] = None
    similarity: float = 0.0
    merge_hint: Optional[str] = None
 
 
class KnowledgeDeduplicator:
    """基于向量相似度的知识去重与合并"""
 
    DUPLICATE_THRESHOLD = 0.92
    PARTIAL_THRESHOLD = 0.75
 
    def __init__(self, chroma_client: chromadb.Client, collection_name: str,
                 embed_fn, llm_client):
        self.collection = chroma_client.get_or_create_collection(
            name=collection_name,
            metadata={"hnsw:space": "cosine"}
        )
        self.embed = embed_fn
        self.llm = llm_client
 
    def _get_embedding_text(self, entry: KnowledgeEntry) -> str:
        """生成用于向量化的文本(标题+摘要+问题模式)"""
        return f"{entry.title}\n{entry.summary}\n{entry.problem_pattern}"
 
    async def check_duplicate(self, entry: KnowledgeEntry) -> DedupDecision:
        """检查新条目是否与已有知识重复"""
        embed_text = self._get_embedding_text(entry)
        vector = self.embed(embed_text)
 
        results = self.collection.query(
            query_embeddings=[vector],
            n_results=5,
            include=["documents", "metadatas", "distances"]
        )
 
        if not results["ids"][0]:
            return DedupDecision(result=DedupResult.NEW)
 
        # distances 是余弦距离(0=相同,1=正交,2=相反),转为相似度
        best_similarity = 1 - results["distances"][0][0]
        best_id = results["ids"][0][0]
 
        logger.info(f"[Dedup] '{entry.title}' 最高相似度: {best_similarity:.3f}")
 
        if best_similarity >= self.DUPLICATE_THRESHOLD:
            return DedupDecision(
                result=DedupResult.DUPLICATE,
                existing_id=best_id,
                similarity=best_similarity
            )
        elif best_similarity >= self.PARTIAL_THRESHOLD:
            # 进一步判断是互补还是矛盾
            existing_doc = results["documents"][0][0]
            conflict_check = await self._check_conflict(entry, existing_doc)
            return DedupDecision(
                result=conflict_check,
                existing_id=best_id,
                similarity=best_similarity
            )
 
        return DedupDecision(result=DedupResult.NEW)
 
    async def _check_conflict(
        self, new_entry: KnowledgeEntry, existing_doc: str
    ) -> DedupResult:
        """用 LLM 判断新旧知识是互补还是矛盾"""
        prompt = f"""判断以下两条知识是互补关系还是矛盾关系。
 
知识A(已有):
{existing_doc[:1000]}
 
知识B(新增):
{new_entry.solution}
 
仅回答:
- "complement":两者互补,可合并
- "conflict":两者矛盾,保留双方并标记冲突
"""
        response = await self.llm.chat(prompt)
        if "conflict" in response.lower():
            return DedupResult.CONFLICT
        return DedupResult.PARTIAL
 
    async def merge_entries(
        self, new_entry: KnowledgeEntry, existing_id: str
    ) -> KnowledgeEntry:
        """合并两条部分重叠的知识条目"""
        existing = self.collection.get(ids=[existing_id], include=["documents", "metadatas"])
        existing_doc = existing["documents"][0] if existing["documents"] else ""
        existing_meta = existing["metadatas"][0] if existing["metadatas"] else {}
 
        prompt = f"""将两条相关知识合并为一条更完整的知识条目。
 
知识A(已有,置信度 {existing_meta.get('confidence', 0.5)}):
{existing_doc[:2000]}
 
知识B(新增,置信度 {new_entry.confidence}):
{new_entry.solution}
 
合并原则:
1. 保留两者都有的信息中置信度更高的版本
2. 互补信息全部保留
3. 合并 tags,去重
4. 置信度取两者最大值
 
输出合并后的 JSON(格式同原条目):"""
 
        response = await self.llm.chat(prompt)
        import re
        match = re.search(r'\{.*\}', response, re.DOTALL)
        if match:
            import json as _json
            data = _json.loads(match.group())
            new_entry.solution = data.get("solution", new_entry.solution)
            new_entry.tags = list(set(new_entry.tags + data.get("tags", [])))
            new_entry.confidence = max(
                new_entry.confidence,
                data.get("confidence", new_entry.confidence)
            )
        return new_entry
 
    async def upsert(self, entry: KnowledgeEntry) -> DedupDecision:
        """去重检查 + 写入/更新"""
        decision = await self.check_duplicate(entry)
 
        if decision.result == DedupResult.DUPLICATE:
            logger.info(f"[Dedup] 跳过重复: '{entry.title}' (相似度 {decision.similarity:.2f})")
            return decision
 
        if decision.result == DedupResult.PARTIAL:
            logger.info(f"[Dedup] 合并部分重叠: '{entry.title}'")
            entry = await self.merge_entries(entry, decision.existing_id)
            # 删除旧条目,写入合并后的新条目
            self.collection.delete(ids=[decision.existing_id])
 
        if decision.result == DedupResult.CONFLICT:
            logger.warning(f"[Dedup] 发现矛盾知识: '{entry.title}',两者均保留")
            entry.title = f"[冲突] {entry.title}"
            entry.tags.append("knowledge-conflict")
 
        # 写入向量库
        embed_text = self._get_embedding_text(entry)
        vector = self.embed(embed_text)
        doc_id = entry.source_fingerprint
 
        self.collection.upsert(
            ids=[doc_id],
            embeddings=[vector],
            documents=[entry.summary + "\n" + entry.solution],
            metadatas={
                "title": entry.title,
                "confidence": str(entry.confidence),
                "tags": ",".join(entry.tags),
                "source_type": entry.source_type,
            }
        )
        logger.info(f"[Dedup] 写入知识库: '{entry.title}'")
        return decision

7. 自动化触发机制

7.1 文件监听(watchdog)

监听 Obsidian Vault 的 inbox 目录,新文件自动触发摄取:

# hermes/triggers/file_watcher.py
import asyncio
from pathlib import Path
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler, FileCreatedEvent
 
 
class InboxWatcher(FileSystemEventHandler):
    """监听 inbox 目录,新文件触发摄取"""
 
    SUPPORTED_EXTENSIONS = {".md", ".txt", ".pdf", ".html"}
 
    def __init__(self, pipeline, inbox_path: str):
        self.pipeline = pipeline
        self.inbox_path = Path(inbox_path)
        self._queue = asyncio.Queue()
 
    def on_created(self, event: FileCreatedEvent):
        if event.is_directory:
            return
        path = Path(event.src_path)
        if path.suffix.lower() in self.SUPPORTED_EXTENSIONS:
            print(f"[Watcher] 发现新文件: {path.name}")
            asyncio.run_coroutine_threadsafe(
                self._queue.put(path),
                asyncio.get_event_loop()
            )
 
    async def run(self):
        observer = Observer()
        observer.schedule(self, str(self.inbox_path), recursive=False)
        observer.start()
        print(f"[Watcher] 监听目录: {self.inbox_path}")
        try:
            while True:
                file_path = await self._queue.get()
                await self.pipeline.ingest_file(file_path)
        finally:
            observer.stop()
            observer.join()

7.2 Claude Code Hooks 集成

~/.claude/settings.json 中配置会话结束钩子,自动触发经验提取:

{
  "hooks": {
    "Stop": [
      {
        "matcher": "",
        "hooks": [
          {
            "type": "command",
            "command": "python3 /path/to/hermes/cli.py extract-session --session-log ~/.claude/session_latest.jsonl"
          }
        ]
      }
    ]
  }
}

对应的 CLI 入口:

# hermes/cli.py
import asyncio
import click
import json
from pathlib import Path
 
 
@click.group()
def cli():
    """Hermes 知识摄取 CLI"""
    pass
 
 
@cli.command()
@click.option("--session-log", required=True, help="会话日志文件路径")
@click.option("--min-length", default=500, help="最短会话长度(字符)")
def extract_session(session_log: str, min_length: int):
    """从 Claude Code 会话日志提取工程经验"""
    log_path = Path(session_log)
    if not log_path.exists():
        click.echo(f"[Error] 会话日志不存在: {log_path}")
        return
 
    # 读取会话内容
    messages = []
    with open(log_path) as f:
        for line in f:
            try:
                messages.append(json.loads(line))
            except json.JSONDecodeError:
                continue
 
    if not messages:
        return
 
    # 重建会话文本
    session_text = "\n".join(
        f"[{m.get('role', 'unknown')}]: {m.get('content', '')}"
        for m in messages
        if m.get("content")
    )
 
    if len(session_text) < min_length:
        click.echo(f"[Skip] 会话太短 ({len(session_text)} chars)")
        return
 
    asyncio.run(_extract_and_save(session_text))
 
 
async def _extract_and_save(session_text: str):
    from hermes.pipeline import HermesPipeline
    from hermes.core import RawContent, SourceType
 
    pipeline = HermesPipeline.from_config("~/.hermes/config.yaml")
    source = RawContent(
        source_type=SourceType.AI_SESSION,
        content=session_text
    )
    entry = await pipeline.ingest(source)
    if entry:
        click.echo(f"[OK] 已摄取: {entry.title} (置信度 {entry.confidence:.0%})")
    else:
        click.echo("[Skip] 未提取到有效知识")
 
 
@cli.command()
@click.argument("url")
def ingest_url(url: str):
    """摄取指定 URL 的网页内容"""
    asyncio.run(_ingest_url(url))
 
 
async def _ingest_url(url: str):
    from hermes.pipeline import HermesPipeline
    from hermes.core import RawContent, SourceType
 
    pipeline = HermesPipeline.from_config("~/.hermes/config.yaml")
    source = RawContent(
        source_type=SourceType.WEB,
        content="",
        source_url=url
    )
    entry = await pipeline.ingest(source)
    if entry:
        click.echo(f"[OK] 摄取成功: {entry.title}")
    else:
        click.echo("[Fail] 摄取失败")
 
 
if __name__ == "__main__":
    cli()

7.3 定时任务配置

# crontab -e 加入以下行
 
# 每天 23:30 处理积压的摄取队列
30 23 * * * python3 /path/to/hermes/cli.py process-queue
 
# 每周一 09:00 生成摄取统计报告
0 9 * * 1 python3 /path/to/hermes/cli.py weekly-report

8. 摄取质量控制

8.1 自动质量评分

每条知识条目摄取后自动打分,低质量条目进入人工审核队列:

# hermes/quality.py
from dataclasses import dataclass
from ..core import KnowledgeEntry
 
 
@dataclass
class QualityScore:
    total: float           # 0.0 - 1.0
    information_density: float
    actionability: float
    specificity: float
    needs_review: bool
 
 
def score_entry(entry: KnowledgeEntry) -> QualityScore:
    """
    三个维度:
    - information_density: 信息密度(字数、是否有具体数据)
    - actionability: 可操作性(解决方案是否具体)
    - specificity: 具体性(是否有代码/命令/数字)
    """
    import re
 
    # 信息密度:solution 字数 + 是否有代码块
    solution_len = len(entry.solution)
    has_code = bool(re.search(r'`[^`]+`|```', entry.solution))
    info_density = min(1.0, solution_len / 300) * 0.7 + (0.3 if has_code else 0)
 
    # 可操作性:solution 是否包含动词/步骤/命令
    action_words = ["执行", "运行", "修改", "设置", "调用", "run", "set", "modify", "call"]
    has_action = any(w in entry.solution.lower() for w in action_words)
    actionability = 0.8 if has_action else 0.4
 
    # 具体性:是否有数字、代码、具体文件路径
    has_numbers = bool(re.search(r'\d+', entry.solution))
    has_path = bool(re.search(r'/\w+|\.py|\.java|\.cpp', entry.solution))
    specificity = min(1.0, (has_numbers + has_path) * 0.5)
 
    total = (
        entry.confidence * 0.4 +
        info_density * 0.3 +
        actionability * 0.2 +
        specificity * 0.1
    )
 
    return QualityScore(
        total=round(total, 2),
        information_density=round(info_density, 2),
        actionability=round(actionability, 2),
        specificity=round(specificity, 2),
        needs_review=total < 0.5 or entry.confidence < 0.5
    )

8.2 摄取统计

# hermes/stats.py
from collections import defaultdict
from datetime import date
import json
from pathlib import Path
 
 
class IngestionStats:
    """摄取统计,持久化到 JSON 文件"""
 
    def __init__(self, stats_file: str = "~/.hermes/stats.json"):
        self.path = Path(stats_file).expanduser()
        self._data = self._load()
 
    def _load(self) -> dict:
        if self.path.exists():
            return json.loads(self.path.read_text())
        return {"daily": {}, "by_source": defaultdict(int), "total": 0}
 
    def record(self, source_type: str, quality_score: float, skipped: bool = False):
        today = date.today().isoformat()
        if today not in self._data["daily"]:
            self._data["daily"][today] = {"ingested": 0, "skipped": 0, "avg_quality": 0}
 
        if skipped:
            self._data["daily"][today]["skipped"] += 1
        else:
            day = self._data["daily"][today]
            n = day["ingested"]
            day["avg_quality"] = (day["avg_quality"] * n + quality_score) / (n + 1)
            day["ingested"] += 1
            self._data["by_source"][source_type] = \
                self._data["by_source"].get(source_type, 0) + 1
            self._data["total"] += 1
 
        self.path.parent.mkdir(parents=True, exist_ok=True)
        self.path.write_text(json.dumps(self._data, ensure_ascii=False, indent=2))
 
    def summary(self) -> str:
        total = self._data["total"]
        by_source = self._data["by_source"]
        last_7_days = list(self._data["daily"].items())[-7:]
        recent_avg = sum(d["avg_quality"] for _, d in last_7_days) / max(len(last_7_days), 1)
 
        return (
            f"知识库统计:\n"
            f"  总条目: {total}\n"
            f"  来源分布: {dict(by_source)}\n"
            f"  近7天平均质量分: {recent_avg:.2f}\n"
        )

9. 完整 Pipeline 代码

9.1 配置文件

# ~/.hermes/config.yaml
llm:
  provider: "anthropic"   # anthropic / openai / ollama
  model: "claude-3-5-haiku-20241022"
  api_key_env: "ANTHROPIC_API_KEY"
  base_url: null          # OpenAI 兼容接口可以设置
 
embedding:
  provider: "local"       # local(sentence-transformers)或 openai
  model: "BAAI/bge-small-zh-v1.5"
 
obsidian:
  vault_path: "~/Documents/ObsidianVault"
  inbox_path: "~/Documents/ObsidianVault/inbox"
  knowledge_path: "~/Documents/ObsidianVault/Knowledge"
  local_rest_api: "http://localhost:27123"
  api_key_env: "OBSIDIAN_API_KEY"
 
vector_db:
  provider: "chroma"
  path: "~/.hermes/chroma_db"
  collection: "knowledge_base"
 
quality:
  min_score: 0.5          # 低于此分进入人工审核队列
  auto_skip_below: 0.2    # 低于此分直接跳过
 
web:
  jina_api_key_env: "JINA_API_KEY"
 
triggers:
  file_watch:
    enabled: true
    inbox_path: "~/Documents/ObsidianVault/inbox"
  webhook:
    enabled: false
    port: 8765

9.2 LLM 客户端统一封装

# hermes/llm_client.py
import os
from abc import ABC, abstractmethod
 
 
class LLMClient(ABC):
    @abstractmethod
    async def chat(self, prompt: str) -> str:
        pass
 
 
class AnthropicClient(LLMClient):
    def __init__(self, model: str, api_key: str = None):
        import anthropic
        self.client = anthropic.AsyncAnthropic(
            api_key=api_key or os.environ["ANTHROPIC_API_KEY"]
        )
        self.model = model
 
    async def chat(self, prompt: str) -> str:
        msg = await self.client.messages.create(
            model=self.model,
            max_tokens=2048,
            messages=[{"role": "user", "content": prompt}]
        )
        return msg.content[0].text
 
 
class OpenAIClient(LLMClient):
    def __init__(self, model: str, api_key: str = None, base_url: str = None):
        from openai import AsyncOpenAI
        self.client = AsyncOpenAI(
            api_key=api_key or os.environ.get("OPENAI_API_KEY", "sk-xxx"),
            base_url=base_url
        )
        self.model = model
 
    async def chat(self, prompt: str) -> str:
        resp = await self.client.chat.completions.create(
            model=self.model,
            messages=[{"role": "user", "content": prompt}],
            max_tokens=2048
        )
        return resp.choices[0].message.content
 
 
class OllamaClient(LLMClient):
    def __init__(self, model: str, base_url: str = "http://localhost:11434"):
        self.model = model
        self.base_url = base_url
 
    async def chat(self, prompt: str) -> str:
        import httpx
        async with httpx.AsyncClient(timeout=120.0) as client:
            resp = await client.post(
                f"{self.base_url}/api/generate",
                json={"model": self.model, "prompt": prompt, "stream": False}
            )
            return resp.json()["response"]
 
 
def create_llm_client(config: dict) -> LLMClient:
    provider = config.get("provider", "anthropic")
    model = config.get("model", "claude-3-5-haiku-20241022")
    api_key_env = config.get("api_key_env")
    api_key = os.environ.get(api_key_env) if api_key_env else None
    base_url = config.get("base_url")
 
    if provider == "anthropic":
        return AnthropicClient(model=model, api_key=api_key)
    elif provider == "openai":
        return OpenAIClient(model=model, api_key=api_key, base_url=base_url)
    elif provider == "ollama":
        return OllamaClient(model=model, base_url=base_url or "http://localhost:11434")
    else:
        raise ValueError(f"不支持的 LLM provider: {provider}")

9.3 Obsidian 写入器

# hermes/writers/obsidian_writer.py
import httpx
from pathlib import Path
from datetime import datetime
from ..core import KnowledgeEntry
 
 
class ObsidianWriter:
    """通过 Obsidian Local REST API 写入知识库"""
 
    def __init__(self, vault_path: str, knowledge_path: str,
                 api_url: str = "http://localhost:27123", api_key: str = None):
        self.vault_path = Path(vault_path).expanduser()
        self.knowledge_path = Path(knowledge_path).expanduser()
        self.api_url = api_url
        self.headers = {"Authorization": f"Bearer {api_key}"} if api_key else {}
 
    def _make_filename(self, entry: KnowledgeEntry) -> str:
        """生成文件名:日期前缀 + 标题"""
        date_str = datetime.now().strftime("%Y%m%d")
        # 清理文件名中的非法字符
        safe_title = "".join(
            c for c in entry.title
            if c.isalnum() or c in "- _()()"
        ).strip()[:50]
        return f"{date_str}-{safe_title}.md"
 
    async def write(self, entry: KnowledgeEntry) -> bool:
        """写入知识条目到 Obsidian"""
        content = entry.to_obsidian_markdown()
        filename = self._make_filename(entry)
        target_path = self.knowledge_path / entry.source_type / filename
 
        # 优先使用 REST API(实时同步)
        if self.api_url:
            try:
                return await self._write_via_api(target_path, content)
            except Exception as e:
                print(f"[ObsidianWriter] REST API 失败,fallback 到文件写入: {e}")
 
        # Fallback:直接写文件
        return self._write_direct(target_path, content)
 
    async def _write_via_api(self, path: Path, content: str) -> bool:
        vault_relative = str(path.relative_to(self.vault_path))
        async with httpx.AsyncClient() as client:
            resp = await client.put(
                f"{self.api_url}/vault/{vault_relative}",
                content=content.encode("utf-8"),
                headers={**self.headers, "Content-Type": "text/markdown"},
                timeout=10.0
            )
            return resp.status_code in (200, 204)
 
    def _write_direct(self, path: Path, content: str) -> bool:
        path.parent.mkdir(parents=True, exist_ok=True)
        path.write_text(content, encoding="utf-8")
        return True

9.4 主 Pipeline

# hermes/pipeline.py
import asyncio
import logging
import os
from pathlib import Path
from typing import Optional
import yaml
 
from .core import RawContent, KnowledgeEntry, SourceType
from .llm_client import create_llm_client
from .ingesters.web_ingester import WebIngester
from .ingesters.session_ingester import AISessionIngester
from .ingesters.generic_ingester import GenericIngester
from .deduplicator import KnowledgeDeduplicator, DedupResult
from .writers.obsidian_writer import ObsidianWriter
from .quality import score_entry
from .stats import IngestionStats
 
logger = logging.getLogger("hermes.pipeline")
 
 
class HermesPipeline:
    """知识摄取主 Pipeline"""
 
    def __init__(self, config: dict):
        self.config = config
 
        # 初始化 LLM
        self.llm = create_llm_client(config["llm"])
 
        # 初始化 Embedding
        self.embed_fn = self._init_embedder(config["embedding"])
 
        # 初始化摄取器链(优先级顺序)
        self.ingesters = [
            WebIngester(self.llm, config.get("web", {})),
            AISessionIngester(self.llm),
            GenericIngester(self.llm),
        ]
 
        # 初始化向量库 + 去重器
        import chromadb
        chroma_path = Path(config["vector_db"]["path"]).expanduser()
        chroma_client = chromadb.PersistentClient(path=str(chroma_path))
        self.deduplicator = KnowledgeDeduplicator(
            chroma_client=chroma_client,
            collection_name=config["vector_db"]["collection"],
            embed_fn=self.embed_fn,
            llm_client=self.llm
        )
 
        # 初始化 Obsidian 写入器
        obs_cfg = config["obsidian"]
        api_key = os.environ.get(obs_cfg.get("api_key_env", ""), "")
        self.obsidian = ObsidianWriter(
            vault_path=obs_cfg["vault_path"],
            knowledge_path=obs_cfg["knowledge_path"],
            api_url=obs_cfg.get("local_rest_api", ""),
            api_key=api_key
        )
 
        self.stats = IngestionStats()
        self.quality_config = config.get("quality", {})
 
    def _init_embedder(self, embed_config: dict):
        provider = embed_config.get("provider", "local")
        model = embed_config.get("model", "BAAI/bge-small-zh-v1.5")
 
        if provider == "local":
            from sentence_transformers import SentenceTransformer
            _model = SentenceTransformer(model)
            return lambda text: _model.encode(text).tolist()
        elif provider == "openai":
            import openai
            client = openai.OpenAI()
            return lambda text: client.embeddings.create(
                input=text, model=model
            ).data[0].embedding
        else:
            raise ValueError(f"不支持的 embedding provider: {provider}")
 
    @classmethod
    def from_config(cls, config_path: str) -> "HermesPipeline":
        path = Path(config_path).expanduser()
        config = yaml.safe_load(path.read_text())
        return cls(config)
 
    async def ingest(self, source: RawContent) -> Optional[KnowledgeEntry]:
        """单条摄取入口"""
        logger.info(f"[Pipeline] 开始摄取 [{source.source_type.value}] {source.source_url or source.fingerprint}")
 
        # 选择摄取器
        ingester = next((i for i in self.ingesters if i.can_handle(source)), None)
        if not ingester:
            logger.warning("[Pipeline] 没有合适的摄取器")
            return None
 
        # 提炼知识条目
        try:
            entry = await ingester.ingest(source)
        except Exception as e:
            logger.error(f"[Pipeline] 摄取失败: {e}", exc_info=True)
            return None
 
        if not entry:
            self.stats.record(source.source_type.value, 0, skipped=True)
            return None
 
        # 质量评分
        quality = score_entry(entry)
        logger.info(f"[Pipeline] 质量分: {quality.total:.2f} | '{entry.title}'")
 
        if quality.total < self.quality_config.get("auto_skip_below", 0.2):
            logger.info(f"[Pipeline] 质量过低,跳过")
            self.stats.record(source.source_type.value, quality.total, skipped=True)
            return None
 
        if quality.needs_review:
            entry.tags.append("needs-review")
            logger.info(f"[Pipeline] 标记为待审核")
 
        # 去重检查 + 写入向量库
        decision = await self.deduplicator.upsert(entry)
        if decision.result.value == DedupResult.DUPLICATE.value:
            self.stats.record(source.source_type.value, quality.total, skipped=True)
            return None
 
        # 写入 Obsidian
        success = await self.obsidian.write(entry)
        if success:
            logger.info(f"[Pipeline] 写入成功: {entry.title}")
        else:
            logger.error(f"[Pipeline] Obsidian 写入失败: {entry.title}")
 
        self.stats.record(source.source_type.value, quality.total)
        return entry
 
    async def ingest_file(self, file_path: Path) -> Optional[KnowledgeEntry]:
        """摄取本地文件"""
        content = file_path.read_text(encoding="utf-8", errors="ignore")
        source_type = SourceType.MANUAL if file_path.suffix == ".md" else SourceType.DOCUMENT
        source = RawContent(
            source_type=source_type,
            content=content,
            source_file=str(file_path)
        )
        return await self.ingest(source)
 
    async def ingest_batch(self, sources: list[RawContent]) -> list[KnowledgeEntry]:
        """批量摄取(并发控制)"""
        sem = asyncio.Semaphore(3)  # 最多同时处理 3 个
 
        async def bounded_ingest(src):
            async with sem:
                return await self.ingest(src)
 
        results = await asyncio.gather(*[bounded_ingest(s) for s in sources])
        return [r for r in results if r is not None]

9.5 快速启动脚本

# run_hermes.py
"""Hermes Pipeline 快速启动脚本"""
import asyncio
import sys
from hermes.pipeline import HermesPipeline
from hermes.core import RawContent, SourceType
 
 
async def main():
    pipeline = HermesPipeline.from_config("~/.hermes/config.yaml")
 
    # 示例 1: 摄取 URL
    if len(sys.argv) > 1 and sys.argv[1].startswith("http"):
        url = sys.argv[1]
        source = RawContent(source_type=SourceType.WEB, content="", source_url=url)
        entry = await pipeline.ingest(source)
        print(f"摄取结果: {entry.title if entry else '失败'}")
 
    # 示例 2: 摄取文本
    elif len(sys.argv) > 1:
        text = open(sys.argv[1]).read()
        source = RawContent(source_type=SourceType.MANUAL, content=text)
        entry = await pipeline.ingest(source)
        print(f"摄取结果: {entry.title if entry else '失败'}")
 
    # 示例 3: 启动文件监听模式
    else:
        from hermes.triggers.file_watcher import InboxWatcher
        watcher = InboxWatcher(pipeline, "~/Documents/ObsidianVault/inbox")
        print("Hermes 已启动,监听 inbox 目录...")
        await watcher.run()
 
 
if __name__ == "__main__":
    asyncio.run(main())

总结

本文构建了一个完整的知识摄取 Pipeline,核心设计要点如下:

维度设计决策
信源覆盖4种类型,按自动化程度分层处理
去重策略内容哈希(完全一致)+ 向量相似度(语义重复)双层保障
质量控制3维度自动评分 + 低分进人工审核队列
幂等性SHA-256 指纹 + 向量库 upsert 保证多次摄取结果一致
LLM 提炼统一 Prompt 模板,支持 Anthropic/OpenAI/Ollama 三种 provider
触发机制watchdog 文件监听 + Claude Code Hooks + crontab 三路触发
存储双写Obsidian Markdown(人类可读)+ Chroma 向量库(机器可检索)

最高价值的摄取来源是 AI 分析会话——每次使用 Claude Code 分析问题后,自动通过 Stop Hook 触发 Hermes,将解决过程提炼为可复用的知识条目,逐步构建起个人工程经验知识库。


文档属于”AI 知识库系统调研”系列,上一篇:05-Obsidian Vault 组织结构设计,下一篇:07-知识检索与 RAG 架构设计。