ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

ZenML 实时事件流(Live Event Streaming)全指南:开启服务端、选配 Broker 与消费 SSE 事件流

ZenML 实时事件流(Live Event Streaming)全指南:开启服务端、选配 Broker 与消费 SSE 事件流 ZenML 实时事件流Live Event Streaming全指南开启服务端、选配 Broker 与消费 SSE 事件流【免费下载链接】zenmlZenML : One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenmlZenML 可以把 pipeline run 内部发布的事件实时推送给任何订阅该 run 的 HTTP 客户端——无论是 LLM token 流式输出、长任务进度更新还是实时 Dashboard 刷新都能在 step 尚未结束时就把中间结果送达前端。本文以仓库文档 docs/book/getting-started/deploying-zenml/live-event-streaming.md 为骨架结合 src/zenml/zen_server/streaming 与 src/zenml/streaming/publishing.py 的源码实现完整讲解服务端如何开启流式传输、如何选择与配置 Broker、如何按 SSE 线协议消费事件流以及投递语义、容量限制与排障手段。读完本文你将能独立完成服务端启用 → 生产端发布 → 客户端消费 → 断线恢复的整条链路搭建与调优。生产端 Python API在 step 内调用zenml.streaming.publish()的完整用法见 Streaming Events本文只在其末尾做速览衔接重点放在服务端操作与线上协议。一、实时事件流解决什么问题传统 MLOps 中step 的输出只有在整个 step 返回后才能以 artifact 或 metadata 的形式被查看中间过程对用户是黑盒。实时事件流Live Event Streaming改变了这一点运行中的 step 可以通过 publishing.py 中的生产者 API 发布事件服务端将其写入 Broker默认实现为 Redis Streams再通过 Server-Sent EventsSSE端点广播给所有订阅者。典型应用场景包括LLM token 流式输出大模型逐 token 生成时实时推送到 Web 前端长任务进度更新数据处理、模型训练等耗时 step 的阶段性进度实时 Dashboard在 run 执行期间动态展示中间指标Agent 工具调用日志把 agent 的决策过程实时同步给观察者。从源码结构看整条链路由四个部分协作完成对应 src/zenml/zen_server/streaming 目录组件源码位置职责生产者_StreamPublishersrc/zenml/streaming/publishing.pystep 内排队事件、后台批量 POST 到服务端Broker 抽象StreamBrokersrc/zenml/zen_server/streaming/brokers/base.py可插拔的流存储/广播后端广播器StreamBroadcastersrc/zenml/zen_server/streaming/broadcaster.py每个 run 一个 reader 会话向订阅者扇出SSE 编码层src/zenml/zen_server/streaming/sse.py把 Broker 条目编码为 SSE 帧处理心跳与过滤⚠️必须提前明确的定位流式传输是 best-effort不是持久化存储。事件有大小上限、负载高时可能被丢弃、Broker 保留窗口过后即消失。一旦事件丢失就永久丢失——没有二级存储、没有回放端点、没有兜底。如果需要留存请在 step 中把关键结果写成 run metadata 或 artifact。二、启用流式传输流式传输默认关闭。唯一的开关是服务端配置项stream_broker_implementation_sourceHelm chart 中对应streaming.streamBrokerImplementationSource。在该字段被设置之前流式端点一律返回501 Not Implemented见 runs_endpoints.py 的streaming_enabled()依赖与 501 响应定义生产端publish()调用被静默丢弃、不发送任何 HTTP 请求服务端不会打开任何 Broker 连接。源码中的开关判断很直接见 server_config.pystreaming_enabled属性即stream_broker_implementation_source is not None。2.1 选择一个 BrokerZenML 内置 Redis Streams Broker完整类名为zenml.zen_server.streaming.brokers.redis_streams.RedisStreamsBroker实现位于 src/zenml/zen_server/streaming/brokers/redis_streams.py。其要求与特性Redis 5.0以及redisPython 扩展pip install zenml[server-streaming]使用 Redis Stream 作为每个 run 的流存储通过XADD MAXLEN ~追加条目transactionFalse的 pipeline 保证热路径开销低每次发布后还会EXPIRE刷新流的 TTL见 redis_streams.py按 deployment ID 命名空间隔离流键stream key 由stream_key_for_run()生成因此多个 ZenML 服务端可以共享同一个 Redis 集群而互不冲突支持单机与 Cluster 模式create_redis_client()会先探测cluster_enabled自动选择Redis或RedisCluster客户端见 redis_client.py。2.2 用 Helm 配置在 Helm values 中设置server: streaming: streamBrokerImplementationSource: zenml.zen_server.streaming.brokers.redis_streams.RedisStreamsBroker environment: ZENML_REDIS_BROKER_URL: redis://my-redis.svc.cluster.local:6379/0Chart 会自动安装一条SSE-only 的 Gateway APIHTTPRoute规则对于携带Accept: text/event-stream且路径位于/api/v1/runs/树下的请求关闭 Envoy 默认的 15 秒请求超时。浏览器的EventSource以及 ZenML 服务端自身发出的帧都满足该条件而自定义客户端如果发送的是 quality-list 形式的Accept头则会落入默认规则在 15 秒后被切断。2.3 用环境变量配置不使用 Helm 部署时通过环境变量设置同样的字段ZENML_SERVER_STREAM_BROKER_IMPLEMENTATION_SOURCEzenml.zen_server.streaming.brokers.redis_streams.RedisStreamsBroker ZENML_REDIS_BROKER_URLredis://...自定义 ingress 的注意点需要在代理层为 SSE 路径关闭请求超时request timeout与响应缓冲response buffering。服务端发出的 SSE 响应已携带以下响应头来覆盖常见中间件见 sse.pyCache-Control: no-cache, no-store, no-transformno-transform防止 gzip 重压缩代理缓冲分块X-Accel-Buffering: no禁用 nginx 风格缓冲但最稳妥的做法是在自己的代理上同样设置这些头。2.4 服务端配置参考以下配置项定义在 server_config.py 中字段ServerConfigurationHelm keyserver.streaming.*默认值说明stream_broker_implementation_sourcestreamBrokerImplementationSource未设置设置此项即启用流式传输。streaming_heartbeat_secondsheartbeatSeconds30.0SSE 心跳间隔源码中Field(default30.0, gt0.0)见 server_config.py。streaming_max_subscribers_per_streammaxSubscribersPerStream100每个 run 的最大并发订阅者数。第 101 个订阅者收到503源码中Field(default100, gt0)且广播器在超过上限时抛出StreamCapacityError见 broadcaster.py。streaming_broadcaster_idle_grace_secondsbroadcasterIdleGraceSeconds30.0最后一个订阅者断开后服务端保持该 run 的 Broker reader 存活的时间这样快速重连无需重新建立 reader见 broadcaster.py 的_delayed_close。2.5 Redis 连接设置连接参数统一从共享的ZENML_REDIS_前缀读取因此同一个 Redis 实例可同时服务于流式 Broker 与其他需要 Redis 的 ZenML 组件。流式 Broker 专属参数使用ZENML_REDIS_STREAMS_BROKER_前缀设置后覆盖共享值前缀合并逻辑见 redis_streams.py 的RedisStreamsBrokerSettings.prefixes()。环境变量默认值说明ZENML_REDIS_BROKER_URL—redis://...或rediss://...。必填。ZENML_REDIS_MAX_CONNECTIONS10连接池大小。并发 run 较多时考虑调大。可用ZENML_REDIS_STREAMS_BROKER_MAX_CONNECTIONS单独覆盖源码约束ge1, le50见 redis_client.py。ZENML_REDIS_SOCKET_TIMEOUT2.0单次调用的 socket 超时秒范围1.0~30.0。ZENML_REDIS_STREAMS_BROKER_MAX_STREAM_LENGTH10000每个 run 保留的最大条目数XADD MAXLEN ~近似裁剪源码默认10_000且ge1。ZENML_REDIS_STREAMS_BROKER_STREAM_TTL_SECONDS3600每个 run 流的 TTL秒每次发布都会EXPIRE刷新。它界定了暂停的生产者多久之内回来不会丢历史。启动健康检查服务端启动时会针对 Broker 做一次连通性检查create_redis_client()中默认ping_on_startTruePING 失败或超时抛出RedisHealthCheckError见 redis_client.py。如果配置的 Redis URL 错误或主机不可达服务端会启动失败并直接报错而不是在后续每个请求上都返回503。三、消费事件流事件流通过 Server-Sent Events (SSE) 暴露在以下端点GET /api/v1/runs/{pipeline_run_id}/events/stream Accept: text/event-stream Authorization: Bearer token权限模型消费需要对该 run 的READ权限与在 Dashboard 中查看该 run 所需权限相同发布则需要UPDATE权限发布端点在 runs_endpoints.py 中通过Action.UPDATE校验。3.1 浏览器端消费const es new EventSource( /api/v1/runs/${runId}/events/stream, { withCredentials: true } ); es.addEventListener(event, (e) console.log(JSON.parse(e.data))); es.addEventListener(end, () es.close());EventSource会自动携带标准的Last-Event-ID请求头进行重连因此短暂的连接中断会从最后收到的事件之后继续见下文断线恢复。3.2 命令行消费curl -N -H Accept: text/event-stream \ -H Authorization: Bearer $ZENML_TOKEN \ $ZENML_URL/api/v1/runs/$RUN_ID/events/stream-N--no-buffer关闭 curl 的输出缓冲让帧在服务端写出的同时立即到达终端。四、SSE 线上格式Wire Format服务端发出的每个 SSE 帧形如id: broker-assigned id event: kind data: JSON-encoded StreamEvent帧由 sse.py 的format_sse_frame()编码id、event、data三个字段按序排列任何字段含换行/回车都会触发校验失败ValueError从而保证线上格式的干净。data字段是 StreamEvent 的 JSON 序列化结果。4.1 保留事件名保留事件名定义在 types.py 的SSEEventName枚举中属于公开线协议不可随意改名event:含义event默认或任意自定义kind生产端发布的事件载荷。data为 JSON 序列化的StreamEvent。endrun 已进入终态服务端将关闭连接。gap订阅者可能错过了最后一个id之后的事件。原因GapReasonoutageBroker 可达性/reader 错误、overflow单订阅者队列已满、shutdown服务端正在关闭。error服务端瞬时错误。客户端应以Last-Event-ID重连。cursor服务端为被过滤掉的事件以及前向兼容的未知帧类型发出的帧。携带id:以推进Last-Event-ID。被过滤事件时data为{}未知帧时data为{unknown_type: type}可用于发现生产端与服务端版本不匹配。客户端可以忽略这两类cursor帧。实现细节sse.py 的_frame_for()Broker 条目若解码为EndFrame则发end帧并终结连接GapMarker编码为gap帧data为{reason: ...}解码失败/未知类型的帧走cursor帧推进游标不满足过滤器的事件同样以cursor帧推进游标保证重连不回放事件的kind若含换行等 SSE 非法字符会丢弃该事件但推进游标避免重连死循环。心跳以 SSE 注释帧: ping\n\n的形式每streaming_heartbeat_seconds默认 30 秒发送一次常量定义于 sse.py发送逻辑见sse_stream()的asyncio.wait超时分支。注释帧不会触发任何addEventListener回调——这正是被过滤/未知帧必须使用event: cursor而非注释帧的原因。4.2 事件过滤SSE 端点接受三个可重复的多值查询参数来限定投递范围。每个参数内多值取 OR参数之间取 AND。被过滤掉的事件依然通过cursor帧推进服务端游标——客户端用Last-Event-ID重连时不会看到它们被重放。参数匹配字段示例kindsStreamEvent.kind?kindstokenkindsprogressstep_namesStreamEvent.step_name即 step 的 invocation id?step_namessummarizecorrelation_idsStreamEvent.correlation_id生产端设置的子流程标签?correlation_idsgen-42组合使用GET /api/v1/runs/{run}/events/stream?kindstokenstep_namessummarize只返回summarizestep 产生的token类型事件。过滤器实现在 sse.py 的EventFilter.matches()任一激活的过滤器不匹配即拒绝。4.3 断线恢复服务端在重连时遵循标准 SSELast-Event-ID请求头语义。浏览器的EventSource会自动发送其他客户端应记录收到的最后一个id:并在重连时带回GET /api/v1/runs/{run}/events/stream Last-Event-ID: last id you received无法设置请求头的客户端某些嵌入式环境可以用?sinceid查询参数作为等价替代——两者都指定起始游标若同时发送请求头优先。恢复语义与边界对应 broadcaster.py 的 catch-up 逻辑订阅者重连后先执行一次追平catch-up从游标起非阻塞读取 Broker 历史再切入实时队列catch-up 与实时之间用 LRU 窗口去重_CATCHUP_IDS_MAX 4096防止边界重叠重复投递若游标早于 Broker 的保留窗口缺失的事件不会重投服务端也不会发信号提示发生了丢失下一次读取返回仍保留的内容订阅者可以挂到已终止的 run服务端回放 Broker 保留的事件历史在保留 TTL 内然后以end事件关闭连接TTL 过期后历史消失订阅直接返回end对应 sse.py 的stale_run_close_response()。丢失事件不可恢复。流式传输是 best-effort事件从不离开 Broker 进入任何持久存储ZenML 也不保留二级副本。Artifact 与 run metadata 持久化的是 run 的结果而不是中间流。请据此设计需要回放能力时在 step 中把关键状态写成 artifact 或 metadata消费者若维护由流派生的 UI 状态累加聚合、滚动缓冲要设计成容忍缺口——丢弃累计状态、基于新事件向前重建而不是指望补取错过的。五、投递语义属性你能得到什么顺序每个 run 内单调递增按 Broker 分配的 id。重复单条连接内每个事件 id 至多投递一次。携带Last-Event-ID重连时服务端严格从最后一个已见 id 之后继续因此不会重投生产端发布失败无重试生产者不会引入重复。订阅者仍建议按事件id做防御性去重。丢失可能的丢失路径生产端队列溢出每进程 4096、服务端发布失败记日志但不重试、Broker 侧MAXLEN裁剪、保留 TTL 过期、单订阅者队列溢出。单订阅者溢出会发出gap: overflow帧其余丢失模式是静默的。丢失事件无法从任何其他来源恢复——ZenML 不保留流的持久副本。保留ZENML_REDIS_STREAMS_BROKER_STREAM_TTL_SECONDS默认最后一次发布后 1 小时。多副本Broker 按 deployment id 键控跨副本投递事件。持久化无。需要持久存储请使用 run metadata 或 artifact。源码印证broadcaster.py每个订阅者一个容量 1024 的asyncio.Queue_SUBSCRIBER_QUEUE_MAXSIZE慢订阅者队列满时丢弃最旧条目并插入GapMarker(reasonOVERFLOW)_put_or_dropreader 以 256 条/批、阻塞 1000ms 的方式读取 Broker_READER_BLOCK_MSBroker 错误时以 0.5s→30s 的指数退避重连带 full jitter并广播限流5 秒内最多一次的gap: outage帧_handle_reader_error/_GAP_RATE_LIMIT_S服务端shutdown时对所有会话广播gap: shutdown与end然后取消 reader 任务shutdown()。六、容量限制单个事件在线上信封内最大64 KiB。常量定义于 constants.pySTREAM_EVENT_PAYLOAD_BYTES_MAX 64 * 1024Broker 帧信封再预留 4 KiB 开销见 frames.py。生产端在publish()内即做编码后大小校验publishing.py超限事件本地直接抛ValueError不会白白占用一次 HTTP 往返。生产端进程内队列最多4096 个事件_QUEUE_MAXSIZE见 publishing.py。队列满时丢弃最旧事件腾出空间publish()本身永不阻塞。每个 run 的 Broker 流默认最多 10000 条XADD MAXLEN ~近似裁剪。订阅者落后过多时会被静默裁剪——线上没有任何针对保留丢失的信号被裁剪的事件也不存在任何可恢复的存储。每个 run 的订阅者上限为streaming_max_subscribers_per_stream默认 100。第 101 个连接收到503 Service Unavailable响应头携带Retry-After: 5广播器侧抛StreamCapacityError路由层转换为 503。生产端批处理参数constants.pypublisher 默认每批最多64 个事件发一次 HTTP 请求可用ZENML_STREAM_PUBLISHER_BATCH_SIZE覆盖注意服务端有批量上限1000超过会导致每次发送都在服务端校验失败。七、故障排查SSE 连接在 ingress 后面 15 秒被切断。你的代理在强制执行请求超时。内置 Helm chart 已为 SSE 配置 Gateway APIHTTPRoute关闭该超时若使用自定义 ingress请对/api/v1/runs/.../events/stream或任何携带Accept: text/event-stream的路径做同样处理。订阅者重连时报告缺失事件。订阅者落后于 Broker 的保留窗口。错过的已永久丢失——没有任何持久存储。解决方案降低生产端速率、调大ZENML_REDIS_STREAMS_BROKER_MAX_STREAM_LENGTH或让订阅者在每个gap帧到达时丢弃累积的流派生状态、基于新事件向前重建。没有事件到达。依次确认流式传输已启用流式端点返回的不是501、消费者对 run 有READ权限、生产端确实在 step 或 pipeline 上下文中调用zenml.streaming.publish()上下文外的调用会被丢弃见 publishing.py。流式端点返回501 Not Implemented。stream_broker_implementation_source未设置。一旦配置完成发布端点与 SSE 端点会同时可用流式传输禁用时它们一起返回501路由层见 runs_endpoints.py 的streaming_enabled()依赖。服务端启动失败报 Stream broker startup probe failed。配置的 Broker 无法触达其后端存储。对 Redis 而言检查ZENML_REDIS_BROKER_URL、TLS 设置以及从服务端 Pod 到 Redis 的网络可达性。启动 PING 失败会抛RedisHealthCheckErrorredis_client.py这是刻意的快速失败——宁可启动时报错也不要在后续每个请求上返回503。发布端点返回503Retry-After: 5。Broker 发布失败runs_endpoints.py事件在服务端被丢弃并记日志不做重试。此时需检查 Redis 连接与MAXLEN/TTL 相关配置。八、生产端 API 速览作为服务端协议的对应面生产端在 step 内的用法详见 Streaming Eventsfrom 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(payload, *, kindevent, correlation_idNone, indexNone)自动从 step 上下文解析 pipeline run 与 step非阻塞入队后台线程批量发送publishing.pykind即线上 SSE 的event:字段客户端用addEventListener(token, ...)订阅end、gap、error、cursor、system是保留名_RESERVED_KINDS生产端直接抛ValueErrorcorrelation_id对 ZenML 透明透传用于给同一逻辑子流程如某次 LLM 生成、某次工具调用的事件分组客户端可在 SSE 端点按?correlation_ids过滤服务端一旦返回501生产者会自静默_disable_publishing()本进程剩余生命周期内所有publish()直接丢弃不发送 HTTP恢复需重启 pipeline 进程需要确保某事件已送达服务端如在发送外部 webhook 前发ready事件时可调用flush(timeout2.0)等待队列排空。九、源码地图与进一步阅读关注点仓库位置服务端配置字段src/zenml/config/server_config.py路由与权限发布/订阅端点src/zenml/zen_server/routers/runs_endpoints.pyBroker 抽象与 Redis 实现src/zenml/zen_server/streaming/brokers/base.py、src/zenml/zen_server/streaming/brokers/redis_streams.pyRedis 客户端与健康检查src/zenml/zen_server/streaming/redis_client.py广播器会话、订阅者扇出、重连src/zenml/zen_server/streaming/broadcaster.pySSE 帧编码与过滤src/zenml/zen_server/streaming/sse.py线协议类型事件名/缺口原因src/zenml/zen_server/streaming/types.pyBroker 帧格式Event/End/Unknownsrc/zenml/zen_server/streaming/brokers/frames.py生产端发布器src/zenml/streaming/publishing.py生产端文档Streaming Events需要特别说明的是流式传输的定位是低延迟的中间过程可视性而非可靠消息队列。在把关键业务事件接入流式通道前务必对照本文投递语义与容量限制两节评估可容忍的丢失窗口需要强持久化的数据请走 run metadata 与 artifact 通道。【免费下载链接】zenmlZenML : One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表