Haystack Pipeline API 深度解析:同步、异步与流式执行的完整指南

发布时间:2026/9/12 1:50:06
Haystack Pipeline API 深度解析:同步、异步与流式执行的完整指南 Haystack Pipeline API 深度解析同步、异步与流式执行的完整指南【免费下载链接】haystackOpen-source AI orchestration framework for building context-engineered, production-ready LLM applications. Design modular pipelines and agent workflows with explicit control over retrieval, routing, memory, and generation. Built for scalable agents, RAG, multimodal applications, semantic search, and conversational systems.项目地址: https://gitcode.com/GitHub_Trending/ha/haystack导读Pipeline是 Haystack 框架的核心编排引擎负责按照执行图调度组件、传递数据并汇总输出。本文以 pipeline_api.md 为骨架结合 pipeline.py 源码系统讲解 Pipeline 的同步执行run、异步执行run_async/run_async_generator、流式输出stream/PipelineStreamHandle以及断点调试与快照恢复break_point/pipeline_snapshot。读完本文你将能根据场景正确选择执行方式、控制并发与输出范围并利用断点机制对复杂管线进行逐步调试与断点续跑。Pipeline 在 Haystack 中的定位在 Haystack 中Pipeline继承自PipelineBase见 base.py是根据执行图运行组件的编排引擎。它通过add_component()注册组件、connect()建立组件间的输入输出连接随后在每次run*调用时按组件名排序保证执行确定性不受插入顺序影响统计每个组件的访问次数component_visits用于判断组件是否可以运行通过优先级队列ComponentPriority反复挑选下一个可运行组件处理READY、BLOCKED、DEFER、HIGHEST等调度状态把叶组件没有出边连接的组件的输出汇总为最终结果。从源码结构看Pipeline同时提供两条执行路径同步路径run()阻塞调用线程直到完成异步路径run_async()、run_async_generator()、stream()。两条路径共享同一套组件注册与图校验逻辑但调度实现不同同步路径在主循环中逐个执行组件见pipeline.py中run的while True循环异步路径则基于asyncio并发调度并用asyncio.Semaphore控制同时运行的组件数量。同步执行Pipeline.run()run()是使用最频繁的入口签名如下run( data: dict[str, Any], include_outputs_from: set[str] | None None, *, break_point: Breakpoint | None None, pipeline_snapshot: PipelineSnapshot | None None, snapshot_callback: SnapshotCallback | None None ) - dict[str, Any]run会同步阻塞调用线程直到管线执行完毕在异步上下文中应改用run_async。完整示例一个最小 RAG 管线文档给出了一个可直接运行的检索增强生成RAG示例——从内存文档库检索、构造提示词、交给 LLM 生成回答from haystack import Pipeline, Document from haystack.components.builders.answer_builder import AnswerBuilder from haystack.components.builders.chat_prompt_builder import ChatPromptBuilder from haystack.components.generators.chat import OpenAIChatGenerator from haystack.components.retrievers.in_memory import InMemoryBM25Retriever from haystack.dataclasses import ChatMessage from haystack.document_stores.in_memory import InMemoryDocumentStore from haystack.utils import Secret # 写入文档到 InMemoryDocumentStore document_store InMemoryDocumentStore() document_store.write_documents([ Document(contentMy name is Jean and I live in Paris.), Document(contentMy name is Mark and I live in Berlin.), Document(contentMy name is Giorgio and I live in Rome.) ]) retriever InMemoryBM25Retriever(document_storedocument_store) prompt_template Given these documents, answer the question. Documents: {% for doc in documents %} {{ doc.content }} {% endfor %} Question: {{question}} Answer: template [ChatMessage.from_user(prompt_template)] prompt_builder ChatPromptBuilder( templatetemplate, required_variables[question, documents], variables[question, documents] ) llm OpenAIChatGenerator() rag_pipeline Pipeline() rag_pipeline.add_component(retriever, retriever) rag_pipeline.add_component(prompt_builder, prompt_builder) rag_pipeline.add_component(llm, llm) rag_pipeline.connect(retriever, prompt_builder.documents) rag_pipeline.connect(prompt_builder, llm) question Who lives in Paris? results rag_pipeline.run( { retriever: {query: question}, prompt_builder: {question: question}, } ) print(results[llm][replies][0].text) # Jean lives in Paris参数详解data管线各组件的输入字典。标准格式是组件名 → 该组件的输入参数即data { comp1: {input1: 1, input2: 2}, }为了方便当输入名在整个管线中唯一时也支持省略组件名、直接平铺data { input1: 1, input2: 2, }data会在内部先经过_prepare_component_input_data()规范化再经validate_input()校验最后转换为内部格式记录每个输入的发送方。从 pipeline.py 可以看到转换后的内部格式形如{component: {socket: [{sender: ..., value: ...}]}}这保证了循环和可变输入variadic场景下数据来源可追溯。include_outputs_from需要额外纳入输出结果集的组件名集合。默认None时返回字典只包含叶组件没有出边连接的组件的输出通过该参数可以把中间组件的输出也放入结果。需要注意对于在循环中被多次调用的组件只包含最后一次产生的输出。break_point断点对象在指定组件运行前触发BreakpointException异常中携带当前管线状态的PipelineSnapshot用于调试。pipeline_snapshot先前中断的管线执行的快照用于断点续跑。可与break_point组合实现步进式调试从快照恢复执行到下一个断点再次暂停。注意break_point必须指向与快照创建时不同的组件或不同的访问次数visit count否则会在恢复后立即再次触发、无法推进。源码在 pipeline.py 中对该冲突做了显式校验并抛出PipelineInvalidPipelineSnapshotError。snapshot_callback创建快照时调用的回调函数。回调接收PipelineSnapshot对象可返回一个可选字符串如文件路径或标识符。提供回调后将取代默认的保存到 JSON 文件行为可用于将快照写入数据库或发送到远程服务。若未提供默认行为是把快照保存为 JSON 文件路径来自_get_output_dir(pipeline_snapshot)见 pipeline.py。返回值与异常返回值dict[str, Any]每个键是组件名值为该组件的输出。include_outputs_from为None时仅含叶组件输出。ValueError向管线提供了非法输入。PipelineRuntimeError管线包含会导致卡死的环unsupported cycles或组件运行失败、返回了不支持的类型。该类定义于 errors.py通过from_exception()/from_invalid_output()工厂方法把组件名、组件类型、错误信息一并封装。PipelineMaxComponentRuns某个组件达到了它在管线中的最大运行次数防止循环失控。PipelineBreakpointException源码中实际类名为BreakpointException见 errors.py断点被触发时抛出包含组件名、状态与部分结果。异步执行Pipeline.run_async()run_async()为管线提供异步接口签名run_async( data: dict[str, Any], include_outputs_from: set[str] | None None, concurrency_limit: int 4, ) - dict[str, Any]concurrency_limit表示允许同时运行的最大组件数默认 4若小于 1 会抛出ValueError。其内部实现实际上是包装run_async_generator——遍历生成器并返回最后一个即最终输出见 pipeline.py。文档中的异步 RAG 示例import asyncio from haystack import Document from haystack.components.builders import ChatPromptBuilder from haystack.components.generators.chat import OpenAIChatGenerator from haystack.components.retrievers.in_memory import InMemoryBM25Retriever from haystack import Pipeline from haystack.dataclasses import ChatMessage from haystack.document_stores.in_memory import InMemoryDocumentStore # 写入文档到 InMemoryDocumentStore document_store InMemoryDocumentStore() document_store.write_documents([ Document(contentMy name is Jean and I live in Paris.), Document(contentMy name is Mark and I live in Berlin.), Document(contentMy name is Giorgio and I live in Rome.) ]) prompt_template [ ChatMessage.from_user( Given these documents, answer the question. Documents: {% for doc in documents %} {{ doc.content }} {% endfor %} Question: {{question}} Answer: ) ] retriever InMemoryBM25Retriever(document_storedocument_store) prompt_builder ChatPromptBuilder(templateprompt_template) llm OpenAIChatGenerator() rag_pipeline Pipeline() rag_pipeline.add_component(retriever, retriever) rag_pipeline.add_component(prompt_builder, prompt_builder) rag_pipeline.add_component(llm, llm) rag_pipeline.connect(retriever, prompt_builder.documents) rag_pipeline.connect(prompt_builder, llm) question Who lives in Paris? async def run_inner(data, include_outputs_from): return await rag_pipeline.run_async(datadata, include_outputs_frominclude_outputs_from) data { retriever: {query: question}, prompt_builder: {question: question}, } results asyncio.run(run_inner(data, include_outputs_from{retriever, llm})) print(results[llm][replies]) # [ChatMessage(_roleChatRole.ASSISTANT: assistant, _content[TextContent(textJean lives in Paris.)], ...)]异步路径的关键实现点run_async_generatorpipeline.py使用asyncio.Semaphore(concurrency_limit)作为并发闸门_schedule_component()为每个 READY 组件创建后台asyncio.Task并在信号量内执行_run_component_async()_run_component_async()对原生异步组件直接await对仅同步组件则通过asyncio.to_thread卸载到线程池执行见 async_utils.py因此同步组件也能混入异步管线具有HIGHEST优先级即含 GreedyVariadic 输入 socket的组件必须单独运行以免其他并发组件向其追加输入造成竞态见_run_component_in_isolation()组件失败时_cancel_in_flight_tasks()会取消并排空仍在运行的任务防止后台泄漏同步组件卸载到的线程无法被中断但其输出会被丢弃不会污染管线状态。逐步产出部分结果Pipeline.run_async_generator()当需要边算边拿如监控每个组件完成情况、尽早展示检索结果时使用run_async_generator( data: dict[str, Any], include_outputs_from: set[str] | None None, concurrency_limit: int 4, ) - AsyncGenerator[dict[str, Any], None]它是一个异步生成器每有一个组件完成就 yield 一份部分输出最终再 yield 完整输出字典。文档示例# 处理结果每完成一个组件即可处理 async def process_results(): async for partial_output in rag_pipeline.run_async_generator( datadata, include_outputs_from{retriever, llm} ): # 每个 partial_output 包含一个已完成的组件的结果 if retriever in partial_output: print(Retrieved documents:, len(partial_output[retriever][documents])) if llm in partial_output: print(Generated answer:, partial_output[llm][replies][0]) asyncio.run(process_results())调度循环在 pipeline.py每轮重建优先级队列 → 取下一个可运行组件 → 尽可能多地并发调度 READY 组件 → 通过_wait_for_tasks(return_whenasyncio.FIRST_COMPLETED)等待任一任务完成并 yield 其输出。若迭代被提前放弃或运行被取消finally分支会取消所有在途任务pipeline.py。异常行为与run_async一致非法输入或concurrency_limit 1抛ValueError组件超限抛PipelineMaxComponentRuns存在不支持的环或组件失败/输出类型非法抛PipelineRuntimeError。实时流式输出Pipeline.stream()与PipelineStreamHandlestream()面向需要逐 token 输出的场景典型如 LLM 打字机效果签名stream( data: dict[str, Any], *, streaming_components: list[str] | None None, include_outputs_from: set[str] | None None, concurrency_limit: int 4, cancel_on_abandon: bool True ) - PipelineStreamHandle使用方式用async for迭代句柄消费StreamingChunk迭代结束后handle.result持有最终管线输出字典与run_async形状相同。默认情况下若消费者中途放弃迭代底层管线任务会被自动取消传cancel_on_abandonFalse则让管线继续跑完。文档中的流式示例import asyncio from haystack.components.builders import ChatPromptBuilder from haystack.components.generators.chat import OpenAIChatGenerator from haystack import Pipeline from haystack.dataclasses import ChatMessage pipe Pipeline() pipe.add_component( prompt_builder, ChatPromptBuilder(template[ChatMessage.from_user(Tell me about {{topic}})]), ) pipe.add_component(llm, OpenAIChatGenerator()) pipe.connect(prompt_builder.prompt, llm.messages) async def main(): handle pipe.stream(data{prompt_builder: {topic: Italy}}) async for chunk in handle: print(chunk.content, end, flushTrue) return handle.result result asyncio.run(main()) print(result[llm][replies])参数与行为streaming_components指定要流式转发的组件名列表。为None默认时转发所有支持流式的组件为列表时只转发列出的组件。传入未知组件名或不支持流式的组件名会抛ValueError。include_outputs_from、concurrency_limit语义与run_async相同。cancel_on_abandon迭代被放弃时是否取消底层管线任务默认True。底层机制pipeline.pystream()先筛选支持异步且暴露streaming_callback输入 socket的组件为每个组件注入一个forwarder回调——把每个StreamingChunk放入内部asyncio.Queue同时若用户提供了streaming_callback初始化时传入或在data中按data{llm: {streaming_callback: cb}}运行时传入也会一并触发。随后创建一个后台任务运行run_async并在finally中放入_END_OF_STREAM哨兵结束队列。推荐使用异步回调同步回调虽被接受但会在事件循环上同步执行可能阻塞事件循环。PipelineStreamHandlePipelineStreamHandle是stream()返回的句柄实现见 pipeline.py对StreamingChunk可异步迭代。result属性最终管线输出字典仅在完整成功运行后才可用。如果管线尚未结束或已被取消访问result会抛RuntimeError若管线运行失败则重新抛出原始异常。源码实现pipeline.py通过检查底层asyncio.Task的状态来区分这三种情况。aclose()方法取消底层管线任务。清理过程受_CLEANUP_TIMEOUT_SECONDS源码中为 1.0 秒限制确保组件无法无限期阻塞清理pipeline.py。迭代语义__aiter__是一个异步生成器async for每次调用都会获得一个生成器退出时执行try/finally——当cancel_on_abandonTrue时放弃迭代会自动取消管线任务为False时任务继续运行至完成。迭代中若管线失败异常会在读取到_END_OF_STREAM哨兵时通过await self._task浮出水面。断点调试与快照恢复run()的三个关键字参数组合起来构成完整的断点 快照调试能力核心数据结构定义在 breakpoints.pyBreakpoint冻结数据类字段为component_name断点所在组件、visit_count组件被访问多少次后才触发默认 0、snapshot_file_path命中时快照保存路径。PipelineState某一时刻的管线状态包括inputs内部格式输入记录每个输入的发送方与到达顺序、component_visits、pipeline_outputs、inputs_format。PipelineSnapshot完整快照含original_input_data、ordered_component_names、pipeline_state、break_point、timestamp、include_outputs_from并校验component_visits与ordered_component_names的一致性。执行机制pipeline.py_run_component/_run_component_async在组件运行前检查断点——若break_point.component_name component_name且visit_count component_visits[component_name]则抛出携带快照的BreakpointException。run()捕获该异常后构造PipelineSnapshot保存组件消费前的输入、原始data、访问计数、已产出的部分结果等将快照挂到异常对象上并通过_save_pipeline_snapshot()落盘为 JSON或交给snapshot_callback自定义处理重新抛出异常由调用方决定如何处理。恢复时run(pipeline_snapshot...)会先校验快照与当前管线图的兼容性然后恢复component_visits、输入状态与已积累的pipeline_outputs从暂停处继续执行。这样便实现了中断 → 检查状态 → 继续执行的完整调试闭环尤其适合 Agent 等复杂多轮管线。选型建议何时用哪种执行方式执行方式适用场景特点run()通用同步调用脚本、同步 Web 框架阻塞线程、确定性执行、支持断点/快照/回调run_async()异步应用如 FastAPI 服务中非阻塞运行基于asyncio并发调度concurrency_limit控并发run_async_generator()需要按组件粒度监控进度、尽早消费中间结果每完成一个组件 yield 一份部分输出stream()LLM 逐 token 打字机输出、长响应边生成边展示返回PipelineStreamHandle可async for消费 chunk四种方式共享同一套组件注册、连接校验与异常体系断点参数目前面向run()同步路径异步路径更侧重于并发与流式。相关实现与测试可继续深入阅读 pipeline.py、errors.py、breakpoints.py以及测试目录下的 test_async_pipeline.py 和 test_breakpoint.py。【免费下载链接】haystackOpen-source AI orchestration framework for building context-engineered, production-ready LLM applications. Design modular pipelines and agent workflows with explicit control over retrieval, routing, memory, and generation. Built for scalable agents, RAG, multimodal applications, semantic search, and conversational systems.项目地址: https://gitcode.com/GitHub_Trending/ha/haystack创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考