跳到主要内容

第 5 章 · 工程内幕:分包、运行时、增量与代码地图

这章讲什么: 前四章讲"算法怎么跑",这章讲"工程怎么搭"——GraphRAG 为什么拆成 8 个包、流水线运行时怎么做到可续跑、存储/缓存/向量库怎么被抽象成可换实现、增量更新怎么只处理新文档。最后给一张贯穿全项目的总代码地图。


5.1 monorepo:按"能力"切成 8 个包

仓库是 uv workspace(pyproject.toml[tool.uv.workspace] members = ["packages/*"])。8 个包各管一件事,主包 graphrag 是"编排壳",其余是可独立复用的能力包:

管什么关键入口
graphrag主壳:workflow 编排、query、CLI、config、data_modelapi/,index/,query/,cli/
graphrag-chunking文本切块策略(token/句子)chunker_factory.create_chunker
graphrag-llmLLM 补全/嵌入、tokenizer、指标、消息构造completion.create_completion,CompletionMessagesBuilder
graphrag-vectors向量库抽象 + 多后端实现vector_store_factory,VectorStore
graphrag-cacheLLM 结果缓存(内存/JSON…)cache_factory,Cache
graphrag-storage输出存储 + 表抽象(文件/blob/cosmos)storage_factory,TableProvider,Table
graphrag-input读入原始文档(csv/txt/json/parquet…)input_reader,TextDocument
graphrag-common跨包公共件(哈希、配置加载、通用工厂基类)hasher,load_config,factory

贯穿全项目的模式:工厂 + 注册表。 PipelineFactory(workflow)、create_chunkercreate_completionvector_store_factorycache_factorystorage_factory 全是同一套路:按配置里的"类型名"造实现,且允许注册自定义实现。带来的性质是每个能力都可替换——想换向量库/换缓存后端/加自定义 workflow,都不用改主流程。


5.2 流水线运行时:共享表 + 状态快照 = 可续跑

第 01 章说过步骤间靠"命名表"交接。这里补运行时的两块骨架(都在 packages/graphrag/graphrag/index/run/run_pipeline.py):

  • PipelineRunContext:一次运行的"公共背包",装着 output_table_provider(读写表)、output_storage(落文件)、cache(LLM 缓存)、callbacks(进度)、state(跨步状态字典)、stats(每步耗时/指标)。每个 workflow 都收到它。
  • 状态快照run_pipeline 启动时会尝试从 output_storage 读回 context.json 恢复 state;每步之后 _dump_stats_json + _dump_context_json 把统计和状态落盘。所以流水线跑到一半失败,重跑能接着来,配合"块 id = 内容哈希"和 LLM 缓存,重跑基本不重复花钱。

出错处理也很干脆:_run_pipeline 用 try/except 包住整个循环,任一步抛异常就 yield 一个带 errorPipelineRunResult 并记 last_workflow,把"哪一步崩的"透出去(见 except 分支)。


5.3 存储与表抽象:同一套代码,多种后端

GraphRAG 要能把中间产物写到本地文件、也能写到 Azure Blob / Cosmos。它用两层抽象隔离:

  • Storagepackages/graphrag-storage/):key-value 式的 blob 读写(file_storage.py/azure_blob_storage.py/memory_storage.py),由 storage_factoryStorageType 选。
  • Table / TableProvidergraphrag-storage/graphrag_storage/tables/):面向"行"的表抽象,支持 async for row in table 流式读、await table.write(row) 流式写、open(name, transformer=...) 打开命名表。前面所有 workflow 的"读表/写表"都走它,所以底层是 parquet 文件还是数据库,workflow 无感。

DataReaderpackages/graphrag/graphrag/data_model/data_reader.py)是查询/报告等步骤读表的便捷门面(reader.entities()reader.relationships()reader.communities()…),把"打开哪张表、转成什么类型"收口在一处。


5.4 LLM 缓存与并发:省钱和限流

索引要调海量 LLM,两件工程事很关键:

  • 缓存graphrag-cache):create_completion(..., cache=..., cache_key_creator=cache_key_creator)cache_key_creatorpackages/graphrag/graphrag/cache/cache_key_creator.py)把"模型+消息+参数"哈希成 key,命中就不再真调 LLM。抽取/合并/报告各自用 context.cache.child(model_instance_name) 拿到隔离的子命名空间,避免串味。
  • 并发:抽取、报告等用 config.concurrent_requests 控制并发线程数、config.async_modeAsyncType.AsyncIO/Threaded)选并发方式;查询期 global 用 asyncio.Semaphore 限并发(第 04 章)。这既压榨吞吐、又不至于把 API 限流打爆。

5.5 增量更新:只处理新来的文档

重新索引整个语料很贵。*-update 方法(第 01 章)在标准流水线尾部接一串 update_* workflow,做"增量合并"。核心是先算 deltapackages/graphrag/graphrag/index/update/incremental_index.py):

# 示意,非源码。对应 get_delta_docs
previous = 已索引文档的标题集合
new_docs = input[~input["title"].isin(previous)] # 新增
deleted_docs = previous 里不在 input# 删除
return InputDelta(new_docs, deleted_docs)

只对 new_docs 跑抽图等重活,再用 concat_dataframes 把新产物接到旧表后面(并重排 human_readable_id)。实体/关系/社区/报告各有对应的 update_*.py 做合并去重。得益于"块 id = 内容哈希",没变的文档天然被识别为"已存在"、不重算。

⚠ 边界:README 提示大版本升级要跑迁移 notebook,且 graphrag init --force 会覆盖配置和 prompt——增量更新是"加文档",不是"跨版本迁移"。


5.6 边界与局限(诚实)

  • 索引贵:官方在 README 顶部就用 ⚠ 警告索引可能很烧钱,建议先小规模试。抽图对每块要多轮 LLM(gleaning)、报告对每个社区再调一次。
  • 质量吃 prompt:默认抽取 prompt 未必贴合你的领域,官方强烈建议先 prompt-tune(第 02 章)。
  • 非官方支持:README 明确说这是"方法演示",不是官方支持的 Microsoft 产品。
  • 全局问题才划算:如果你的问题全是局部事实类,basic 就够了,那份昂贵的图+报告索引未必回本(见 index.md §4 横向对比)。
  • NLP 快速路径有取舍fast 省钱,但共现图没有类型/描述、噪声更多,要靠 prune_graph 补救。

5.7 总代码地图(跨全项目)

主题文件路径符号名
CLI 命令入口packages/graphrag/graphrag/cli/main.py@app.command("init"/"index"/"update"/"query"/"prompt-tune")
索引 APIpackages/graphrag/graphrag/api/index.pybuild_index
查询 APIpackages/graphrag/graphrag/api/query.pyglobal_search,local_search,drift_search,basic_search
流水线工厂packages/graphrag/graphrag/index/workflows/factory.pyPipelineFactory
流水线运行时 + 快照packages/graphrag/graphrag/index/run/run_pipeline.pyrun_pipeline,PipelineRunContext
全局配置模型packages/graphrag/graphrag/config/models/graph_rag_config.pyGraphRagConfig
各步配置模型packages/graphrag/graphrag/config/models/*.pyExtractGraphConfig,GlobalSearchConfig
数据模型 + 表列packages/graphrag/graphrag/data_model/Entity,Relationship,Community,schemas.py
读表门面packages/graphrag/graphrag/data_model/data_reader.pyDataReader
增量 deltapackages/graphrag/graphrag/index/update/incremental_index.pyget_delta_docs,concat_dataframes
LLM 缓存 keypackages/graphrag/graphrag/cache/cache_key_creator.pycache_key_creator
LLM 补全/嵌入包packages/graphrag-llm/graphrag_llm/create_completion,CompletionMessagesBuilder
向量库抽象packages/graphrag-vectors/graphrag_vectors/vector_store_factory,VectorStore
存储/表抽象packages/graphrag-storage/graphrag_storage/storage_factory,TableProvider,Table