
R2R 任务编排深度指南基于 Hatchet 的分布式工作流与知识图谱构建实战【免费下载链接】R2RSoTA production-ready AI retrieval system. Agentic Retrieval-Augmented Generation (RAG) with a RESTful API.项目地址: https://gitcode.com/GitHub_Trending/r2/R2RR2RSoTA production-ready AI retrieval system将文档摄入ingestion与知识图谱knowledge graph构建这类复杂、耗时且可能失败的任务交给 Hatchet 这一分布式、容错的任务队列来编排。本文以官方编排指南 docs/cookbooks/orchestration.md 为骨架结合仓库源码与 Docker 编排配置讲解 Hatchet 的核心概念、R2R 内置的九大工作流、并发与超时参数以及如何通过 7274 端口的图形界面观察、检查甚至重试任务。读完本文你将能解释一次r2r documents create-samples背后发生了什么并具备独立配置、观测与排障 R2R 编排链路的能力。为什么 R2R 需要一层任务编排一条 R2R 摄入流水线远不止上传文件这么简单文件解析parsing、文本分块chunking、向量化embedding、写入存储之后还可能触发文档摘要、上下文增强chunk enrichment乃至知识图谱抽取而知识图谱构建又包含实体关系抽取、实体描述、去重、社区聚类与社区摘要等多个阶段。任何一个阶段失败、任何一个环节耗时过长都会拖垮整体吞吐。R2R 的解法是引入 Hatchet目前支持hatchet与simple两种实现——前者是面向生产的分布式方案后者是默认的进程内直调方案。Hatchet 的三个核心概念官方文档给出了理解 Hatchet 的三块基石Workflows工作流由一组函数step组成的执行单元响应外部触发器如 API 请求而运行R2R 中的摄入一个文件构建一个图都对应一个 workflow。Workers工作者长期运行、负责实际执行 workflow 函数的进程R2R 的 worker 名为r2r-worker见下文工厂代码。Managed Queue托管队列用于处理实时任务、低延迟的队列保证任务提交后能被快速调度。在此基础上Hatchet 还提供 step 依赖parents、并发限制concurrency、重试retries、超时timeout与失败回调on_failure_step等机制这些都被 R2R 的编排实现直接使用。R2R 中的编排架构与配置入口配置模型OrchestrationConfigR2R 的编排配置定义在 py/core/base/providers/orchestration.pyprovider必填仅支持hatchet与simplevalidate_config会校验max_runs 2_048单个 worker 最大并发运行数ingestion_concurrency_limit 16摄入类工作流的并发上限graph_search_results_creation_concurrency_limit 32图谱创建类并发上限graph_search_results_concurrency_limit 8图谱搜索/摘要类并发上限。在 py/core/main/config.py 的REQUIRED_KEYS中orchestration一节强制要求provider键存在默认配置文件 py/r2r/r2r.toml 中[orchestration]的默认值为provider simple即不引入额外基础设施。若你的部署使用 Docker Compose 全套编排见下文应将其改为hatchet以启用分布式队列。工厂装配worker 是如何被创建的在 py/core/main/assembly/factory.py 的create_orchestration_provider中可以看到装配逻辑当config.provider hatchet时创建HatchetOrchestrationProvider并立即get_worker(r2r-worker)否则回退到SimpleOrchestrationProvider。HatchetOrchestrationProviderpy/core/providers/orchestration/hatchet.py是对hatchet-sdk的薄封装workflow()/step()/failure()/concurrency()直接透传 SDK 装饰器run_workflow()通过orchestrator.admin.run_workflow把任务投递到队列并返回task_idregister_workflows()则按工作流类型Workflow.INGESTION/Workflow.GRAPH加载对应的工厂函数并把产出的 workflow 注册到 worker 上。编排带来的三大收益官方文档概括为三点可扩展性Scalability高效处理大规模任务并发上限可配置、worker 可横向扩展容错性Fault Tolerance内置重试机制与错误处理失败时有on_failure回调兜底灵活性Flexibility随着 R2R 能力扩张新增或修改工作流非常容易——新增一个 workflow 只需在工厂函数中注册。R2R 内置工作流全景从文档清单到源码实证官方文档列出五个高层工作流IngestFilesWorkflow、UpdateFilesWorkflow、KgExtractAndStoreWorkflow、CreateGraphWorkflow、EnrichGraphWorkflow。对照仓库源码它们在 Hatchet 中实际注册为九个子工作流分布在两个工厂文件里文档中的概念工作流Hatchet 中的实际工作流name源码位置IngestFilesWorkflowingest-files、ingest-chunksingestion_workflow.pyUpdateFilesWorkflowupdate-chunk、create-vector-index、delete-vector-index同上KgExtractAndStoreWorkflowgraph-extraction含抽取 实体描述两步、graph-deduplicationgraph_workflow.pyCreateGraphWorkflow / EnrichGraphWorkflowgraph-clustering、graph-community-summarization同上摄入工作流ingest-files / ingest-chunksingest-files是核心摄入工作流ingestion_workflow.py其parsestep 内部依次驱动文档状态机PARSING → AUGMENTING生成摘要→ EMBEDDING → STORING → SUCCESS若开启了 chunk enrichment 还会经历ENRICHING。每一步都调用IngestionService的对应方法parse_file产出 extractions 并累计total_tokens摘要阶段调用augment_document_infoembed_document批量生成向量store_embeddings写入向量库finalize_ingestion收尾并把文档/分块关联到 collection同时将集合的graph_sync_status、graph_cluster_status标记为OUTDATED从而触发后续图谱流程。值得注意的源码细节并发控制ingest-files使用orchestration_provider.concurrency(max_runsingestion_concurrency_limit, limit_strategyConcurrencyLimitStrategy.GROUP_ROUND_ROBIN)并发键取自请求中的用户 ID——即同一用户的多个文档串行处理不同用户之间轮询公平调度超时与重试parsestep 设置timeout60m、retries0摄入失败不做盲目重试级联触发若配置了automatic_extraction摄入完成后通过context.aio.spawn_workflow(graph-extraction, ...)自动派生子工作流把知识图谱抽取串进摄入链路失败兜底on_failure回调会把文档状态更新为FAILED并写入step_run_errors()作为失败元数据避免任务静默消失。ingest-chunks则演示了 Hatchet 的多 step 依赖写法ingest解析分块→embedparents[ingest]嵌入并存储→finalizeparents[embed]收尾step 之间通过context.step_output(...)传递中间结果。create-vector-index与delete-vector-index分别负责向量索引的创建timeout 360m与删除timeout 30m。图谱工作流extraction → clustering → community summarization图谱侧的编排在 graph_workflow.py 中timeout360m且关键 step 带retries1充分体现知识图谱构建长耗时、可重试的特点graph-extraction第一步graph_search_results_extraction从文档/集合中抽取实体与关系并入库同时把extraction_status置为PROCESSING若只给collection_id工作流会先列出该集合下所有文档再为每个文档spawn_workflow并发子任务。第二步graph_search_results_entity_descriptionparents[...]为实体生成描述若配置了automatic_deduplication还会级联触发graph-deduplication子工作流graph-clustering第一步执行社区聚类若graph_cluster_status已是SUCCESS会抛出需先 reset graph的R2RException避免重复构建第二步graph_search_results_community_summary按每批最多 100 个社区parallel_communities min(100, num_communities)把社区摘要任务切分为 N 个graph-community-summarization子工作流并行派发全部完成后将extraction_status升级为ENRICHED、graph_cluster_status置为SUCCESSgraph-community-summarization / graph-deduplication分别独立承担单个社区摘要与文档实体去重均带并发限制与失败回调。快速上手通过 Docker 部署 Hatchet 编排栈官方文档指出R2R 的 Docker 发行版默认附带 Hatchet 前端应用端口为7274。在 docker/compose.full.yaml 中可以看到完整编排栈hatchet-postgresHatchet 的元数据库PostgreSQLhatchet-rabbitmq消息队列宿主机端口映射5673:5672、15673:15672hatchet-create-db/hatchet-migration先建库脚本见 docker/scripts/create-hatchet-db.sh再跑 Hatchet 迁移镜像版本hatchet-migrate:v0.53.15hatchet-setup-config执行hatchet-admin quickstart生成配置hatchet-engine引擎服务端口7077:7077hatchet-dashboard前端面板映射7274:80setup-token通过 docker/scripts/setup-token.sh 为固定租户707d0855-80ab-4e1f-a156-f1c4546cbf52创建 API token写入/hatchet_api_key/api_key.txtR2R 服务启动脚本 docker/scripts/start-r2r.sh 会读取该 token 作为HATCHET_CLIENT_TOKEN环境变量供hatchet-sdk连接队列使用。部署完成后浏览器打开http://localhost:7274即可进入 Hatchet 前端使用默认账号登录邮箱Emailadminexample.com密码PasswordAdmin123!!注意以上账号与租户 ID 为当前仓库 Docker 编排中的默认值见 compose.full.yaml 与 setup-token.sh生产环境务必替换。观察与排障Workflow Runs 面板实战触发一次真实任务在 R2R 服务就绪后执行r2r documents create-samples该命令会投递一批示例文档的摄入任务。打开http://localhost:7274/workflow-runs即可看到 Hatchet 工作流面板实时展示每个 workflow run 的状态排队中、运行中、成功、失败这正是官方文档截图所展示的Running workflows场景。检查单个工作流与 GUI 重试点击任意 workflow run 即可进入详情页查看每个 step 的输入输出、执行耗时与日志。官方文档特别强调当任务失败时你可以直接在 GUI 中重试retry该任务无需重新提交整个请求。这在偶发的 LLM 超时、临时网络抖动导致图谱抽取失败时尤为实用——重试前on_failure回调已经先把相关状态如extraction_status标记为失败重试后状态会重新进入PROCESSING。长时任务与超时设计知识图谱构建是典型的长时间运行任务Hatchet 对此有原生支持。从源码可以量化看到 R2R 的超时设计摄入工作流ingest-files/ingest-chunkstimeout60m图谱系列工作流graph-extraction、graph-clustering、graph-community-summarization、graph-deduplication及create-vector-indextimeout360m6 小时worker 级超时同样设置为 60m配合retries1支撑图谱构建这类长任务。官方文档中Worker timeout is set to 60m to support long running tasks like graph construction正是对这一设计的说明——如果你要处理超大文档集或极慢的 LLM可据此合理调整 step 超时与并发上限。从源码验证工作流注册与失败处理链路最后把文档描述与实现串成一条完整的调用链请求进入 R2R 服务后由服务层调用run_workflow(ingest-files, ...)hatchet.py任务进入 Hatchet 托管队列r2r-worker拉取执行工作流定义来自hatchet_ingestion_factory与hatchet_graph_search_results_factory两个工厂函数ingestion_workflow.py、graph_workflow.py它们通过orchestration_provider.workflow(...)装饰器向 worker 注册失败时各工作流的on_failure回调把文档/图谱状态更新为FAILED同时 GUI 保留现场可供重试若使用轻量部署默认配置provider simple则 simple.py 直接异步调用服务方法逻辑等价但无队列与重试能力。这套文档描述概念 → 配置定义参数 → 工厂注册实现 → GUI 观测排障的完整闭环正是 R2R 编排层的设计精髓。随着项目迭代官方文档预告还会继续补充编排特性集与最佳实践——但仅凭当前仓库你已经可以独立完成编排链路的配置、触发、监控与故障恢复。【免费下载链接】R2RSoTA production-ready AI retrieval system. Agentic Retrieval-Augmented Generation (RAG) with a RESTful API.项目地址: https://gitcode.com/GitHub_Trending/r2/R2R创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考