DEP(轻量级):用于 Sweeper 作业进度和候选结果的事件模式和事件层传输
作者: devivasudevan创建于 2026年9月18日更新于 2026年9月20日
标签enhancementlanguage::rustlanguage::pythondynamo-runtimedeployment::natsplannerdep:draftDGDR
目前,正在运行的 Sweeper 搜索是一个黑盒。唯一的可见性来自临时的 tqdm 控制台输出。最终的产品(index.json 和 DGD YAML 文件)仅在搜索完成后才写入。
flowchart LR
A[Sweeper.run starts] --> B[Search loop: rounds of candidates]
B --> C{on_round callback}
C -->|synchronous, unguarded| B
B --> D[Search completes]
D --> E[index.json + DGD YAML written]
E --> F[Publisher can finally read results]
B -.no external visibility.-> X((?))
style X fill:#f66,stroke:#900,color:#fff
style F fill:#9f6,stroke:#090这导致了两个关键缺陷:
- **发布者没有消费机制。**他们只能轮询终端作业状态或在完成后监视文件系统路径 - 没有办法在发现候选项时做出反应,或者在执行长时间搜索期间逐步更新
DGDRRun状态。 - **失败或卡住的搜索在终止之前不会提供任何诊断信号。**对于持续数小时的搜索,消费者无法区分“仍在搜索,已找到 N 个候选项”和“卡住” - 没有事件流来区分这两种情况。
真正的
Sweeper.run签名包括一个on_round回调:on_round: Callable[[int, list[Candidate]], None]。它接收每轮的完整Candidate列表,并在主搜索循环中以同步和未加锁的方式调用。 **在发射中执行的工作会直接停止搜索 - 这是整个设计中的一个硬性约束。
Dynamo 已经运行了一个在生产环境中经过验证的机制,用于将运行进程中的结构化事件流式传输给消费者,而无需轮询 - EventPublisher/EventSubscriber - 目前用于前向传递指标、KV 缓存事件和序列跟踪。
事件模式
一个通用的封装包装了每个事件:
{
"apiVersion": "sweeper.dynamo.nvidia.com/v1alpha1",
"runUID": "3f2b1a90-...",
"sequence": 42,
"timestamp": "2026-09-17T12:00:00Z",
"type": "round_completed",
"data": { }
}apiVersion- 包装的版本信息,以便消费者可以拒绝不熟悉的格式runUID- 标识 Sweeper 运行,用于与DGDRRun进行关联sequence- 每个运行的单调递增值,用于排序type- 以下之一: …
内容来源: ai-dynamo/dynamo