尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

PocketFlow 实战:基于 100 行级 LLM 框架实现流式响应与用户中断(LLM Streaming Interruption)

PocketFlow 实战:基于 100 行级 LLM 框架实现流式响应与用户中断(LLM Streaming  Interruption) PocketFlow 实战基于 100 行级 LLM 框架实现流式响应与用户中断LLM Streaming Interruption【免费下载链接】PocketFlowPocket Flow: 100-line LLM framework. Let Agents build Agents!项目地址: https://gitcode.com/gh_mirrors/poc/PocketFlow导读本文围绕 PocketFlow 仓库中 cookbook/pocketflow-llm-streaming 这一实战示例系统讲解如何利用 PocketFlow 的 Node/Flow 抽象在命令行终端中实现 LLM 响应的实时逐字流式输出以及通过独立监听线程实现随时按 ENTER 中断生成的交互能力。读完本文你将掌握streamTrue的 OpenAI 流式调用接入方式、无 API Key 时的可运行假数据流fake stream模拟方案以及把阻塞式输入监听与流式消费循环解耦的线程化中断模式并能将同样的中断思路复用到你自己的 PocketFlow 流水线中。示例概览什么是 LLM Streaming and Interruption流式输出Streaming是现代 LLM 应用中最常见的体验之一模型不是一次性吐出完整回答而是把 token 一个个吐出来前端随之逐字渲染从而显著降低首字延迟time to first token。在此基础上允许用户在中途打断生成则是流式交互场景的进阶需求——例如用户看到回答方向不对可以立刻停止浪费时间和 token 的生成过程。cookbook/pocketflow-llm-streaming这个示例把这两个需求做成了一个最小可运行的 PocketFlow 应用其核心能力是实时显示LLM 生成的每个内容分片chunk在到达后立即打印无需等待完整响应随时中断在任何时刻按下 ENTER 键即可中断当前流式生成线程随即被干净地收尾。示例的依赖极其精简仅两个文件分别承担流程编排与LLM 调用的职责main.pyStreamNode节点实现包含中断监听线程、流式消费循环与收尾逻辑utils.py真实的 OpenAI 流式调用函数stream_llm与无需密钥即可运行的假流函数fake_stream_llm。快速运行示例的requirements.txt只有两个依赖见 requirements.txtpocketflow0.0.1 openai1.0.0在仓库根目录下按以下步骤运行pip install -r cookbook/pocketflow-llm-streaming/requirements.txt python cookbook/pocketflow-llm-streaming/main.py默认情况下示例使用fake streaming responses假流式响应即无需配置任何 API Key 即可完整体验逐字打印 ENTER 中断的效果。启动后终端会先提示Press ENTER at any time to interrupt streaming...随后文字开始逐块刷新输出此时任意时刻按一下 ENTER程序会打印User interrupted streaming.并结束最后通过flow.run(shared)返回defaultaction。核心实现拆解StreamNode 的 prep - exec - post整个示例只有 49 行代码主体是一个继承自 PocketFlowNode的StreamNode见 main.py。理解它之前先回顾 PocketFlow 框架最基础的约定每个 Node 都包含三个阶段prep(shared) - exec(prep_res) - post(shared, prep_res, exec_res)其中prep负责从共享存储shared中读取和预处理数据exec只做计算不访问sharedpost负责回写结果并返回决定下一步走向的 action 字符串默认default。这一点在框架源码 pocketflow/init.py 与官方文档 docs/core_abstraction/node.md 中有完整定义。StreamNode正是按照这一三段式结构组织的README 中给出的工作流程也对应这三个阶段见 README.md 的 How It Works1. prep创建中断监听线程准备流式数据源def prep(self, shared): interrupt_event threading.Event() def wait_for_interrupt(): input(Press ENTER at any time to interrupt streaming...\n) interrupt_event.set() listener_thread threading.Thread(targetwait_for_interrupt) listener_thread.start() prompt shared[prompt] chunks stream_llm(prompt) return chunks, interrupt_event, listener_thread关键点有三个threading.Event作为线程间协作信号监听线程在用户按下 ENTER 后调用event.set()主线程的消费循环通过event.is_set()感知中断请求。这是跨线程通信最轻量、最安全的方式避免了在消费者线程中直接操作input()。shared共享存储传参Prompt 通过 PocketFlow 的shared字典传入本例为shared {prompt: Whats the meaning of life?}这正是框架推荐的数据流动方式——prep从shared读post向shared写。流式数据源在 prep 中建立stream_llm(prompt)返回的是 OpenAI 的流式响应迭代器而非完整文本。注意中断监听线程在prep阶段就启动因此从流式输出的第一秒起用户就具备打断能力。2. exec消费内容分片并实时渲染def exec(self, prep_res): chunks, interrupt_event, listener_thread prep_res for chunk in chunks: if interrupt_event.is_set(): print(User interrupted streaming.) break if hasattr(chunk.choices[0].delta, content) and chunk.choices[0].delta.content is not None: chunk_content chunk.choices[0].delta.content print(chunk_content, end, flushTrue) time.sleep(0.1) # simulate latency return interrupt_event, listener_thread这一阶段实现了 README 所述的逐块实时显示与处理用户中断每轮循环先检查interrupt_event.is_set()一旦用户已按 ENTER 就立即break实现中断对每个 chunk 使用hasattr(...)is not None双重校验delta.content这是 OpenAI 流式协议的标准健壮性写法——流式响应中可能包含空 delta 或仅携带role字段的起始 chunk直接取.content会抛AttributeErrorprint(..., end, flushTrue)不换行且强制刷新缓冲区是终端逐字打字机效果的关键utils.py的独立运行测试也复用了同样的打印手法见 utils.pytime.sleep(0.1)用于模拟真实网络延迟让流式效果肉眼可见真实调用时可移除。这里体现了 PocketFlow 设计哲学中exec只做计算、不触碰shared的约束流式消费、延迟模拟、中断判定全部发生在纯计算层shared仅在prep/post中被读写。3. post收尾清理并返回 actiondef post(self, shared, prep_res, exec_res): interrupt_event, listener_thread exec_res interrupt_event.set() listener_thread.join() return defaultpost承担两个职责清理监听线程先interrupt_event.set()确保即使流式自然结束阻塞在input()上的监听线程也能尽快返回再用listener_thread.join()等待其结束避免出现悬挂线程或程序无法退出的问题返回 action返回default字符串。本示例只有一个节点Flow 没有注册后续节点因此flow.run()到此结束若注册了后续节点则会按 action 继续流转。4. 组装与运行node StreamNode() flow Flow(startnode) shared {prompt: Whats the meaning of life?} flow.run(shared)Flow(startnode)指定入口节点flow.run(shared)从起始节点开始执行。从框架源码看Flow._orch会在每个节点执行后调用get_next_node根据 action 查找后继节点并继续循环见 pocketflow/init.py本例无后继节点流在StreamNode处自然终止。框架的行为一致性由 tests/test_flow_basic.py 中的用例覆盖验证例如test_start_method_initialization断言了无后继时flow.run()返回最后一个节点post()的返回值。两种流式数据源stream_llm 与 fake_stream_llmutils.py提供了两种可互换的流式数据源这也是 README API Key 一节的核心内容。真实流式调用 stream_llmdef stream_llm(prompt): client OpenAI(api_keyos.environ.get(OPENAI_API_KEY, your-api-key)) response client.chat.completions.create( modelgpt-4o, messages[{role: user, content: prompt}], temperature0.7, streamTrue # Enable streaming ) return response要点streamTrue是流式的开关开启后create()立即返回一个可迭代的响应对象每次迭代产出一个 SSE chunk而不是等待完整回答模型与采样参数默认使用gpt-4o、temperature0.7可按需替换模型名或调整参数流式请求同样支持max_tokens、top_p等 OpenAI 标准参数API Key 来源优先读取环境变量OPENAI_API_KEY未设置时回退到占位字符串your-api-key此时真实调用必然报鉴权错误因此才需要配合 fake 数据源使用。免密钥假流 fake_stream_llmdef fake_stream_llm(prompt, predefined_textThis is a fake response. ...): chunk_size 10 class SimpleObject: def __init__(self, **kwargs): for key, value in kwargs.items(): setattr(self, key, value) for i in range(0, len(predefined_text), chunk_size): text_chunk predefined_text[i:ichunk_size] delta SimpleObject(contenttext_chunk) choice SimpleObject(deltadelta) chunk SimpleObject(choices[choice]) chunks.append(chunk) return chunks它用最简单的动态对象构造出与 OpenAI 流式响应同构的嵌套结构chunk.choices[0].delta.contentSimpleObject通过setattr动态挂载属性结构等价于choices - delta - content并把一段预设文本按chunk_size10切成小块。因此fake_stream_llm返回的对象可以直接喂给StreamNode.exec里同一套hasattr(chunk.choices[0].delta, content)判空逻辑无需改动任何消费代码——这是该示例默认零成本可运行的关键设计。接入真实 OpenAI 流式响应README 给出了从 fake 切换到真实的完整步骤见 README.md第一步编辑 main.py把调用函数从假流替换为真流# Change this line: chunks fake_stream_llm(prompt) # To this: chunks stream_llm(prompt)第二步确保环境变量中已设置 OpenAI API Keyexport OPENAI_API_KEYyour-api-key-here随后再次运行python cookbook/pocketflow-llm-streaming/main.py即可看到来自真实gpt-4o的逐 token 流式输出并同样支持 ENTER 中断。注意几点前提与限制stream_llm内部的OpenAI(api_key...)会覆盖环境变量未设置时的默认占位值因此在没有密钥的情况下必须保持使用fake_stream_llm真实流式的 chunk 到达间隔取决于网络与模型推理速度可考虑移除exec中的time.sleep(0.1)若遭遇限流rate limit或配额错误可参考 PocketFlow Node 内置的max_retries/wait重试机制见 docs/core_abstraction/node.md 的 Fault Tolerance Retries 一节源码实现位于 pocketflow/init.pytests/test_fall_back.py 对其重试与回退行为做了完整验证。原理纵深为什么用独立线程 Event实现中断把等待用户按键放进主消费循环里是不可行的——input()是阻塞调用会卡住流式输出本身导致无法边输出边监听。示例采用的独立监听线程 threading.Event信号是处理这类问题的经典模式监听线程只做一件事阻塞在input()上用户按键后event.set()返回线程使命完成主线程流式消费循环只做一件事遍历 chunks、渲染内容并在每轮循环里以非阻塞方式查询event.is_set()收尾阶段post中set()join()确保监听线程不会残留。这种解耦保证了渲染与监听互不阻塞且中断响应延迟不超过一个 chunk 的消费时间本例加上0.1s模拟延迟后最大响应约 0.1 秒。扩展思路从 Node 到多节点 Flow当前示例是单节点 Flow。基于 PocketFlow 的 action 机制node_a - action node_b详见 docs/core_abstraction/flow.md你可以很自然地把StreamNode融入更大的流水线例如stream_node summary_node # 流式结束后再做摘要 # 或 stream_node - interrupted fallback_node # 被中断时走降级分支只需让StreamNode.post根据interrupt_event.is_set()返回不同 action 字符串即可如中断返回interrupted自然结束返回default这正是 PocketFlow 分支流转的推荐用法其行为语义同样在 tests/test_flow_basic.py 的分支用例中得到验证。文件速览与进一步阅读文件作用cookbook/pocketflow-llm-streaming/main.pyStreamNode节点实现与 Flow 组装完整演示中断线程、流式消费与清理cookbook/pocketflow-llm-streaming/utils.pystream_llm真实流式调用与fake_stream_llm假流模拟含独立自测入口cookbook/pocketflow-llm-streaming/requirements.txt依赖声明pocketflow与openaicookbook/pocketflow-llm-streaming/README.md本示例的官方说明文档pocketflow/init.pyPocketFlow 框架核心源码Node、Flow、重试/回退与 action 流转实现docs/core_abstraction/node.mdNode 三段式抽象prep/exec/post与容错重试官方文档docs/core_abstraction/flow.mdFlow action 驱动流转、分支、循环与嵌套 Flow 官方文档tests/test_flow_basic.py / tests/test_fall_back.py框架流转与重试回退行为的单元测试可用于印证本文所述机制如果你对流式相关能力做进一步探索仓库中还有异步流式与 WebSocket 流式推送的独立示例cookbook/pocketflow-fastapi-websocket可结合本文的终端版流式中断实现一起阅读形成从 CLI 到 Web 的完整流式交互方案。【免费下载链接】PocketFlowPocket Flow: 100-line LLM framework. Let Agents build Agents!项目地址: https://gitcode.com/gh_mirrors/poc/PocketFlow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表