ZenML 实时事件流(Live Event Streaming)全指南:开启服务端、选配 Broker 与消费 SSE 事件流
【免费下载链接】zenmlZenML 🙏: One AI Platform from Pipelines to Agents. https://zenml.io.项目地址: https://gitcode.com/GitHub_Trending/ze/zenml
ZenML 可以把 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 Events(SSE)端点广播给所有订阅者。
典型应用场景包括:
- LLM token 流式输出:大模型逐 token 生成时实时推送到 Web 前端;
- 长任务进度更新:数据处理、模型训练等耗时 step 的阶段性进度;
- 实时 Dashboard:在 run 执行期间动态展示中间指标;
- Agent 工具调用日志:把 agent 的决策过程实时同步给观察者。
从源码结构看,整条链路由四个部分协作完成(对应 src/zenml/zen_server/streaming 目录):
| 组件 | 源码位置 | 职责 |
|---|---|---|
生产者_StreamPublisher | src/zenml/streaming/publishing.py | step 内排队事件、后台批量 POST 到服务端 |
Broker 抽象StreamBroker | src/zenml/zen_server/streaming/brokers/base.py | 可插拔的流存储/广播后端 |
广播器StreamBroadcaster | src/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_source(Helm chart 中对应streaming.streamBrokerImplementationSource)。在该字段被设置之前:
- 流式端点一律返回
501 Not Implemented(见 runs_endpoints.py 的streaming_enabled()依赖与 501 响应定义); - 生产端
publish()调用被静默丢弃、不发送任何 HTTP 请求; - 服务端不会打开任何 Broker 连接。
源码中的开关判断很直接,见 server_config.py:streaming_enabled属性即stream_broker_implementation_source is not None。
2.1 选择一个 Broker
ZenML 内置 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 ~追加条目(transaction=False的 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_SOURCE=zenml.zen_server.streaming.brokers.redis_streams.RedisStreamsBroker ZENML_REDIS_BROKER_URL=redis://...自定义 ingress 的注意点:需要在代理层为 SSE 路径关闭请求超时(request timeout)与响应缓冲(response buffering)。服务端发出的 SSE 响应已携带以下响应头来覆盖常见中间件(见 sse.py):
Cache-Control: no-cache, no-store, no-transform(no-transform防止 gzip 重压缩代理缓冲分块)X-Accel-Buffering: no(禁用 nginx 风格缓冲)
但最稳妥的做法是在自己的代理上同样设置这些头。
2.4 服务端配置参考
以下配置项定义在 server_config.py 中:
字段(ServerConfiguration) | Helm key(server.streaming.*) | 默认值 | 说明 |
|---|---|---|---|
stream_broker_implementation_source | streamBrokerImplementationSource | 未设置 | 设置此项即启用流式传输。 |
streaming_heartbeat_seconds | heartbeatSeconds | 30.0 | SSE 心跳间隔(源码中Field(default=30.0, gt=0.0),见 server_config.py)。 |
streaming_max_subscribers_per_stream | maxSubscribersPerStream | 100 | 每个 run 的最大并发订阅者数。第 101 个订阅者收到503(源码中Field(default=100, gt=0)且广播器在超过上限时抛出StreamCapacityError,见 broadcaster.py)。 |
streaming_broadcaster_idle_grace_seconds | broadcasterIdleGraceSeconds | 30.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_CONNECTIONS | 10 | 连接池大小。并发 run 较多时考虑调大。可用ZENML_REDIS_STREAMS_BROKER_MAX_CONNECTIONS单独覆盖(源码约束ge=1, le=50,见 redis_client.py)。 |
ZENML_REDIS_SOCKET_TIMEOUT | 2.0 | 单次调用的 socket 超时(秒),范围1.0~30.0。 |
ZENML_REDIS_STREAMS_BROKER_MAX_STREAM_LENGTH | 10000 | 每个 run 保留的最大条目数(XADD MAXLEN ~近似裁剪),源码默认10_000且ge=1。 |
ZENML_REDIS_STREAMS_BROKER_STREAM_TTL_SECONDS | 3600 | 每个 run 流的 TTL(秒),每次发布都会EXPIRE刷新。它界定了"暂停的生产者多久之内回来不会丢历史"。 |
启动健康检查:服务端启动时会针对 Broker 做一次连通性检查(create_redis_client()中默认ping_on_start=True,PING 失败或超时抛出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。 |
end | run 已进入终态,服务端将关闭连接。 |
gap | 订阅者可能错过了最后一个id之后的事件。原因(GapReason):outage(Broker 可达性/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重连时不会看到它们被重放。
| 参数 | 匹配字段 | 示例 |
|---|---|---|
kinds | StreamEvent.kind | ?kinds=token&kinds=progress |
step_names | StreamEvent.step_name(即 step 的 invocation id) | ?step_names=summarize |
correlation_ids | StreamEvent.correlation_id(生产端设置的子流程标签) | ?correlation_ids=gen-42 |
组合使用:
GET /api/v1/runs/{run}/events/stream?kinds=token&step_names=summarize只返回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>无法设置请求头的客户端(某些嵌入式环境)可以用?since=<id>查询参数作为等价替代——两者都指定起始游标;若同时发送,请求头优先。
恢复语义与边界(对应 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(reason=OVERFLOW)(_put_or_drop); - reader 以 256 条/批、阻塞 1000ms 的方式读取 Broker(
_READER_BLOCK_MS),Broker 错误时以 0.5s→30s 的指数退避重连(带 full jitter),并广播限流(5 秒内最多一次)的gap: outage帧(_handle_reader_error/_GAP_RATE_LIMIT_S); - 服务端
shutdown时对所有会话广播gap: shutdown与end,然后取消 reader 任务(shutdown())。
六、容量限制
- 单个事件在线上信封内最大64 KiB。常量定义于 constants.py:
STREAM_EVENT_PAYLOAD_BYTES_MAX = 64 * 1024;Broker 帧信封再预留 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.py):publisher 默认每批最多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 失败会抛RedisHealthCheckError(redis_client.py),这是刻意的快速失败——宁可启动时报错,也不要在后续每个请求上返回503。
发布端点返回503+Retry-After: 5。Broker 发布失败(runs_endpoints.py):事件在服务端被丢弃并记日志,不做重试。此时需检查 Redis 连接与MAXLEN/TTL 相关配置。
八、生产端 API 速览
作为服务端协议的对应面,生产端在 step 内的用法(详见 Streaming Events):
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": f"working on item {i}"}) publish({"phase": "done"}) return "ok"关键点:
publish(payload, *, kind="event", correlation_id=None, index=None)自动从 step 上下文解析 pipeline run 与 step,非阻塞入队,后台线程批量发送(publishing.py);kind即线上 SSE 的event:字段,客户端用addEventListener("token", ...)订阅;end、gap、error、cursor、system是保留名(_RESERVED_KINDS),生产端直接抛ValueError;correlation_id对 ZenML 透明透传,用于给同一逻辑子流程(如某次 LLM 生成、某次工具调用)的事件分组,客户端可在 SSE 端点按?correlation_ids=过滤;- 服务端一旦返回
501,生产者会自静默(_disable_publishing()):本进程剩余生命周期内所有publish()直接丢弃,不发送 HTTP(恢复需重启 pipeline 进程); - 需要确保某事件已送达服务端(如在发送外部 webhook 前发"ready"事件)时,可调用
flush(timeout=2.0)等待队列排空。
九、源码地图与进一步阅读
| 关注点 | 仓库位置 |
|---|---|
| 服务端配置字段 | src/zenml/config/server_config.py |
| 路由与权限(发布/订阅端点) | src/zenml/zen_server/routers/runs_endpoints.py |
| Broker 抽象与 Redis 实现 | src/zenml/zen_server/streaming/brokers/base.py、src/zenml/zen_server/streaming/brokers/redis_streams.py |
| Redis 客户端与健康检查 | src/zenml/zen_server/streaming/redis_client.py |
| 广播器(会话、订阅者扇出、重连) | src/zenml/zen_server/streaming/broadcaster.py |
| SSE 帧编码与过滤 | src/zenml/zen_server/streaming/sse.py |
| 线协议类型(事件名/缺口原因) | src/zenml/zen_server/streaming/types.py |
| Broker 帧格式(Event/End/Unknown) | src/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),仅供参考