ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

ZenML 实时事件流(Streaming Events)实战指南:从 Step 内发布事件到客户端订阅

ZenML 实时事件流(Streaming Events)实战指南:从 Step 内发布事件到客户端订阅 ZenML 实时事件流Streaming Events实战指南从 Step 内发布事件到客户端订阅【免费下载链接】zenmlZenML : One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml本指南围绕 ZenML 的实时事件流Live Event Streaming功能展开聚焦生产者端的 Python API如何在step内部调用zenml.streaming.publish()把 LLM 令牌流、长任务进度、实时仪表盘数据等中间输出在 step 返回之前推送给订阅方。读完本文你将掌握publish()/flush()的完整用法、kind与correlation_id的分流与分组技巧、动态 Pipeline 中的事件归属规则以及生产者侧的队列、批量、限流等底层实现细节并了解如何与服务端的 SSE 订阅端点配合完成端到端的事件链路。为什么需要流式事件ZenML Pipeline 中一个 step 的返回值通常要到整个 step 执行完毕才会产生。对于 LLM 文本生成、长时间数据清洗、模型评估这类任务用户往往希望尽早看到中间结果生成到一半的 token、处理到第几条记录、当前处于哪个阶段。ZenML 的实时事件流正是为此设计Pipeline 可以在 step 运行过程中把事件实时发送到服务器任何订阅了该 run 的客户端浏览器、CLI、脚本、仪表盘都能立即收到。它适用于LLM 令牌流逐 token 推送生成内容长任务进度step 运行中的阶段与进度更新实时仪表盘边运行边渲染中间状态任何需要在 step 返回前浮出中间输出的场景。本文对应的是官方文档 Streaming events生产者端 API的完整解读与源码级扩展。服务端如何开启流、选择 broker、订阅 SSE 流及线协议约定参见 Live event streaming本文末尾也会给出与消费端的衔接要点。前置条件重要流式功能默认关闭必须在服务端显式开启设置stream_broker_implementation_source配置项。不调用zenml.streaming.publish()的 Pipeline 完全不受影响。一旦服务器返回501 Not Implemented生产者会在当前进程剩余生命周期内自我静默self-mute后续publish()调用直接返回、不再发送任何 HTTP 请求。从 Step 内部发布事件生产者的核心用法极其简单——在step装饰的函数体内直接调用publish()from zenml import step from zenml.streaming import publish step def my_streaming_step() - str: publish({phase: warmup}) for i in range(10): publish({i: i, msg: fworking on item {i}}) publish({phase: done}) return ok关键点publish()会从 step 上下文step context自动读取当前 Pipeline run 和 step 的信息因此在step函数内部无需传递任何句柄该调用不会阻塞事件先进入进程内的队列由后台线程分批small batches发送到服务器事件是与 step 返回值相互独立的旁路数据——即使 step 最终返回ok订阅方早已实时看到了 warmup → 逐项处理 → done 的完整过程。上下文解析publish()如何知道发给谁从源码看src/zenml/streaming/publishing.pypublish()内部通过_resolve_publish_context()解析归属优先尝试get_step_context()拿到pipeline_run.id、step_run.id和step_name若不在 step 上下文则检查DynamicPipelineRunContext动态 Pipeline 的 run 上下文两者都不存在时返回None——此时事件被丢弃仅记录一条debug级别的日志详见下文动态 Pipeline 与上下文外调用。publish()API 详解from zenml.streaming import publish, flush publish( payload: Dict[str, Any], *, kind: str event, correlation_id: Optional[str] None, index: Optional[int] None, ) - None各参数说明参数类型说明payloadDict[str, Any]任意 JSON 可序列化的字典。单个事件的线协议封套wire envelope上限为64 KiB大小检查在publish()内部同步执行超限会在本地直接抛错而不是等到 HTTP 往返之后才发现。kindstr自由格式的事件类型标签客户端可据此过滤。默认event。correlation_idOptional[str]对 ZenML 不透明不做校验、不解释原样透传给消费者用于给同一逻辑子流程的事件分组或排序例如一次 LLM 生成对应一个 ID。消费者可在 SSE 端点按correlation_id过滤。indexOptional[int]同一correlation_id组内的可选有序序号同样原样透传。参数背后的实现细节结合 src/zenml/streaming/publishing.py 与 src/zenml/models/v2/core/stream_event.pypublish()会校验kind是否撞上保留名end、gap、error、cursor、system是服务端控制帧保留的事件名见 src/zenml/zen_server/streaming/types.py 中的SSEEventName。生产者使用这些kind会直接抛出ValueError避免与线协议的控制帧混淆。从源码可见生产者端保留集合_RESERVED_KINDS与SSEEventName的值一一对应。解析发布上下文上文已述构造StreamEvent模型包含pipeline_run_id、step_run_id、step_name、kind、correlation_id、index、payload七个字段入队交给进程内的单例发布器_StreamPublisher。payload的序列化与大小检查由_check_payload_size()完成publishing.py先通过 pydantic 的TypeAdapter做 JSON 编码编码失败不可序列化或编码后字节数超过STREAM_EVENT_PAYLOAD_BYTES_MAX64 * 1024定义于 src/zenml/constants.py都会抛出ValueError。选择合适的kindkind在线上就是 SSE 帧的event:字段——客户端用addEventListener(token, ...)订阅对应类型。维护一组稳定、可枚举的类型名比用临时标签更容易让消费者长期编码依赖。社区常见选择tokenLLM 流式生成progressstep 进度status状态变化log自由格式的日志行。用correlation_id分组当单个 step 同时为多个并行子流程发事件时例如一个 Agent 并发发起多次 LLM 生成用correlation_id给事件打标消费者即可按逻辑子流程区分事件流publish({text: chunk}, kindtoken, correlation_idgen-42, indexi)对应的消费者订阅时可以带上查询参数GET /api/v1/runs/{run}/events/stream?correlation_idsgen-42只接收这一路生成的 token。correlation_id对 ZenML 而言是完全不透明的——唯一的契约是消费者可以按等值过滤。flush()确认事件已送达from zenml.streaming import flush published flush(timeout2.0) # True 表示队列已排空后台工作线程会持续排空队列并且进程退出时有atexit兜底见 publishing.py首个publish()触发线程启动时注册atexit.register(self.shutdown)所以大多数场景不需要手动调用flush()。何时需要当你在做某件必须在事件到达服务器之后才执行的操作时例如先发布一个ready事件紧接着发送外部 webhook——而该 webhook 的目标正是要消费这条事件流。此时用flush()保证先决事件已送达。flush(timeout)的实现publishing.py会等待缓冲队列排空且当前在途批次发送完成发布器通过_cond条件变量与_sending标志保证flush()不会在队列为空但批次仍在发送中时误报成功。返回True表示在超时前排空超时返回False。动态 Pipeline 与上下文外调用动态 Pipelinepipeline(dynamicTrue)体内在任意step之外调用publish()事件会归属到该 Pipeline run但step_run_id与step_name为空见 publishing.py走DynamicPipelineRunContext分支。详见 Dynamic Pipelines。任何 Pipeline / step 上下文之外模块顶层代码、REPL、独立脚本publish()调用被丢弃仅记录一条debug级别日志——因为没有可归属的 run。官方文档中也明确指出此类调用是dropped after a debug-level log line。生产者侧限制与调优限制项数值说明进程内队列容量4 096 个事件队满时丢弃最旧事件以腾出空间——publish()本身从不阻塞。定义于 publishing.py 的_QUEUE_MAXSIZE。单个事件 payload64 KiB线协议封套超限在本地抛ValueError见上文。每 HTTP 请求批量64 个事件默认可通过ZENML_STREAM_PUBLISHER_BATCH_SIZE覆盖。批量上限须 ≤1 000服务端批量上限STREAM_EVENT_MAX_BATCH_SIZE见 constants.py超限会导致每次发送都被服务端校验拒绝。发布权限UPDATE发布事件需要对 run 拥有UPDATE权限服务端端点实现见下文。底层单例发布器与后台线程_StreamPublisher是一个单例SingletonMetaClass其工作模型如下publishing.pypublish(event)先检查_disabled标志若服务端已报告流式关闭则直接返回再做 payload 大小检查首次调用时启动名为zenml-stream-publisher的守护线程事件入deque缓冲区队满即popleft()丢弃最旧工作线程_run()循环_collect_batch()每次最多取STREAM_PUBLISHER_BATCH_SIZE个事件并原子地置_sending True_send_batch()按pipeline_run_id分组defaultdict每组通过zen_store.publish_run_events()发一次StreamBatchRequest批量模型见 src/zenml/models/v2/core/stream_event.py若服务端返回NotImplementedError或MethodNotAllowedError即流式未开启、端点返回 501调用_disable_publishing()——在当前进程剩余生命周期内静默并打印 warning 日志要恢复需重启 Pipeline 进程其他发送异常仅记录 warning服务端不重试最后_done_sending()清标志并唤醒等待中的flush()。这条链路与官方文档events are queued and a background thread sends them to the server in small batches的描述完全一致。服务端侧对照批量端点与权限校验与服务端对应的接收端点是POST /api/v1/runs/{pipeline_run_id}/events见 src/zenml/zen_server/routers/runs_endpoints.py端点依赖streaming_enabled检查——流式未开启时返回501 Not Implemented进入后校验对 run 的Action.UPDATE权限空批次直接返回count0批次先编码校验再交给stream_broker().publish()写入 brokerbroker 写入失败返回503 Service Unavailable且带Retry-After: 5。消费端快速衔接本文聚焦生产者端完整的订阅协议在 Live event streaming 中。为便于读者串联整条链路这里给出最小衔接服务端开启后订阅端点为GET /api/v1/runs/{pipeline_run_id}/events/stream携带Accept: text/event-stream与 Bearer Token消费需要 run 的READ权限与在 Dashboard 查看该 run 所需权限相同每个 SSE 帧形如id/event/data三段event:即生产者传入的kinddata为 JSON 序列化的StreamEvent客户端可用EventSource浏览器或curl -N命令行订阅kinds、step_names、correlation_ids三个查询参数可做 AND/OR 组合过滤断线后浏览器自动携带Last-Event-ID续传其他客户端可用?sinceid等价恢复。关键取舍尽力而为Best-effort而非持久化使用实时事件流前必须明确其语义边界事件会丢失生产者队列溢出、服务端发布失败仅记录日志、不重试、broker 的MAXLEN截断、保留窗口TTL到期、订阅者队列溢出等都可能丢事件没有二级副本ZenML 不为事件流保留任何持久化拷贝也没有回放端点。一旦事件丢失就永久丢失持久化需求请走正式通道需要保留的内容请写成 run metadata 或 artifact——它们持久化 run 的结果而流式只负责实时传递过程。小结ZenML 的实时事件流在生产者侧被刻意设计得极简一个不阻塞的publish()加上可选的flush()就能让运行中的 step 实时广播中间状态。理解其底层队列模型4 096 事件/进程、批量发送默认 64/请求、上限 1 000、64 KiB payload 上限以及尽力而为的投递语义能帮助你在 LLM 令牌流、长任务进度、实时仪表盘等场景中做出正确设计——把要保留下来的写进元数据与制品把要实时看的送进事件流。延伸阅读服务端开启与运维Live event streaming生产者端实现源码src/zenml/streaming/publishing.py、src/zenml/streaming/init.py事件模型与批量请求模型src/zenml/models/v2/core/stream_event.py服务端 SSE 保留事件名与缺口原因src/zenml/zen_server/streaming/types.py常量与限制定义src/zenml/constants.py服务端发布端点实现src/zenml/zen_server/routers/runs_endpoints.py动态 PipelineDynamic Pipelines持久化替代方案run metadata 与 artifacts【免费下载链接】zenmlZenML : One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表