Dozzle 通知系统架构深度解析:从 Manager 表达式引擎到多主机配置广播

📅 发布时间:2026/9/14 11:32:12
Dozzle 通知系统架构深度解析:从 Manager 表达式引擎到多主机配置广播
Dozzle 通知系统架构深度解析从 Manager 表达式引擎到多主机配置广播【免费下载链接】dozzleRealtime log viewer for containers. Supports Docker, Swarm and K8s.项目地址: https://gitcode.com/GitHub_Trending/do/dozzleDozzle 的告警Alerts子系统负责监视容器日志、资源指标与生命周期事件并在用户定义的表达式条件命中时把通知投递到 Webhook 或 Dozzle Cloud。本文基于仓库中的 通知系统架构记忆文档结合 Manager 实现、订阅类型定义 与 多主机服务 等源码完整拆解该子系统的核心组件、事件处理流水线、类型映射的坑以及并发模型帮助读者理解 Dozzle 告警从规则编译到通知发送的全链路并掌握排查该子系统常见缺陷的思路。架构总览Manager 是订阅与调度的中枢通知系统的核心是notification.Manager其结构定义在 manager.go 中// Manager manages notification subscriptions and dispatches notifications type Manager struct { subscriptions *xsync.Map[int, *Subscription] dispatchers *xsync.Map[int, dispatcher.Dispatcher] cloudDispatcher atomic.Pointer[dispatcher.Dispatcher] subscriptionCounter atomic.Int32 dispatcherCounter atomic.Int32 listener *ContainerLogListener statsListener *ContainerStatsListener eventListener *ContainerEventListener ctx context.Context cancel context.CancelFunc sendSem *semaphore.Weighted }从源码结构看Manager 的职责边界非常清晰subscriptions / dispatchers两张xsync.Map分别持有订阅规则与投递器Webhook、Cloud所有并发读写都走并发安全的 mapcloudDispatcher用atomic.Pointer单独存放云投递器订阅的DispatcherID 0时被解释为投递到 Cloud见getDispatcher三个监听器日志、指标、事件只负责把原始数据整理成带容器/主机元信息的通道真正的规则匹配与发送都收敛在 Manager 内sendSem是一个容量为 5 的semaphore.WeightedNewManager中semaphore.NewWeighted(5)用于限制并发通知发送数防止慢速 Webhook 拖垮整个系统。NewManager在构造时立即启动事件处理协程这是该子系统的启动关键点// Start processing log events from the listener go m.processLogEvents() // Start processing stat events from the stats listener go m.processStatEvents() // Start processing Docker events from the event listener go m.processDockerEvents()记忆文档中提到两个处理协程processLogEvents与processStatEvents对照当前源码NewManager实际还启动了processDockerEvents来处理容器生命周期事件对应事件类告警三者一一对应 Dozzle 的三类告警日志告警、指标告警、事件告警。Manager 还实现了容器匹配接口ShouldListenToContainer只有存在已启用且带日志表达式的订阅匹配该容器时日志监听器才会为它建立日志流。纯指标订阅不会触发日志流——这是一个重要的资源优化点避免为只关心 CPU 的容器白白拉取全量日志。三类监听器的生命周期差异日志监听器与指标监听器在生命周期管理上并不对称这是理解整个子系统行为的关键日志监听器一旦启动常开Manager.Start()调用listener.Start(m)后日志监听器即进入常开状态。它用activeStreams *xsync.Map[string, *streamEntry]containerID → 活跃日志流见 log_listener.go跟踪每个容器的流订阅变化时通过UpdateStreams增删流而不是整体重启。指标监听器随订阅启停指标监听器ContainerStatsListener定义了完整的 start/stop 生命周期见 stats_listener.goStart()订阅所有 client 的 stats 通道启动enrich协程Stop()通过取消 context 取消全部订阅IsRunning()供外部查询状态。Manager 在每次订阅增删改后都会调用updateListeners()做联动决策if hasMetric { m.statsListener.Start() } else { m.statsListener.Stop() } if hasEvent { m.eventListener.Start() } else { m.eventListener.Stop() }也就是说只要没有任何启用的指标告警stats 订阅会被整体停掉SubscribeStats不会消耗任何资源。事件监听器同理。enrich循环本身也值得细看它用带缓冲的chan *ContainerStatEvent容量 1000承载富化后的事件并用一个5 秒 TTL 缓存避免每秒的 stats tick 都去重查容器与主机信息——缓存注释解释了原因挂载盘剩余空间由卷监控器带外刷新5 秒的 TTL 保证mounts[*].usedPercent这类表达式能在几秒内看到新值。当通道满时事件直接丢弃并告警日志Metric stats channel full, dropping stat event而不是阻塞上游。记忆文档曾标记过stats_listener 的enrich()不检查 channel 是否关闭可能空转的风险。对照当前源码enrich的接收分支已包含if !ok { return }的关闭检查stats_listener.go该问题在现版本中已不成立——这类记忆 → 源码的逐条核对正是维护该子系统时最有价值的做法。数据模型订阅、表达式与冷却语义Subscription 结构订阅定义在 types.go 中同时承载可持久化字段与纯运行时字段均以json:- yaml:-标记不落盘type Subscription struct { ID int json:id yaml:id Name string json:name yaml:name Enabled bool json:enabled yaml:enabled DispatcherID int json:dispatcherId yaml:dispatcherId LogExpression string json:logExpression yaml:logExpression ContainerExpression string json:containerExpression yaml:containerExpression MetricExpression string json:metricExpression,omitempty yaml:metricExpression,omitempty EventExpression string json:eventExpression,omitempty yaml:eventExpression,omitempty Cooldown int json:cooldown,omitempty yaml:cooldown,omitempty // 秒指标通知间隔 SampleWindow int json:sampleWindow,omitempty yaml:sampleWindow,omitempty // 评估窗口秒数 // 运行时状态不持久化 TriggerCount atomic.Int64 // 触发次数 LastTriggeredAt atomic.Pointer[time.Time] // 最近触发时间 TriggeredContainerIDs *xsync.Map[string, struct{}] // 触发过的容器去重集合 MetricCooldowns *xsync.Map[string, time.Time] // 每容器指标冷却 EventCooldowns *xsync.Map[string, time.Time] // 每容器事件冷却 MetricSampleBuffers *xsync.Map[string, *utils.RingBuffer[bool]] // 每容器采样环 }四条表达式分别在CompileExpressions()中用 expr-lang 编译成*vm.Program且各自绑定不同的求值环境types.NotificationContainer、types.NotificationLog、types.NotificationStat、types.NotificationEvent。这意味着表达式可用的字段由对应环境的结构体 tag 决定例如日志表达式里message可能是字符串也可能是 map复杂日志匹配失败类型不匹配时按不匹配处理并打 Debug 日志而不是报错。冷却与采样窗口的取值语义记忆文档提到前端NotificationRule.cooldown是可选的后端默认 300。对照源码需要修正一个细节300 是前端默认值指标告警表单中const cooldown ref(props.alert?.cooldown ?? ... ?? 300)见 MetricAlertFields.vue后端GetCooldownSeconds()并不注入 300它只做区间裁剪见 types.go// GetCooldownSeconds returns the cooldown in seconds, clamped to [0, 3600] func (s *Subscription) GetCooldownSeconds() int { if s.Cooldown 0 { return 0 } if s.Cooldown 3600 { return 3600 } return s.Cooldown }即 0 表示无冷却上限被钳制到 3600 秒。采样窗口GetSampleWindowSeconds()同理默认 15 秒、上限 300 秒RecordMetricSample用每容器的RingBuffer[bool]记录窗口内每次求值结果窗口填满且 ≥80% 样本命中才触发告警——这是防抖的核心机制单点毛刺不会误报。类型映射的坑跨包转换时的字段完整性记忆文档中Type Mapping Gotchas一节是整个通知子系统最值得警惕的部分。Dozzle 的包结构里存在三套相似但不同的配置类型notification.Subscription/notification.DispatcherConfig内部运行态、types.SubscriptionConfig/types.DispatcherConfig对 Agent 广播用、以及 Cloud 配置。任何一处手工字段拷贝都可能漏字段。container.ContainerStat 与 types.NotificationStat 的字段差异指标表达式求值用的是types.NotificationStat见 types/notification.gotype NotificationStat struct { CPUPercent float64 json:cpu expr:cpu MemoryPercent float64 json:memory expr:memory MemoryUsage float64 json:memoryUsage expr:memoryUsage Mounts []NotificationMount json:mounts,omitempty expr:mounts }而容器侧原始统计container.ContainerStat见 internal/container/types.go字段更多还包含NetworkRxTotal、DiskWriteTotal等。两个结构体的字段名相同Go 字段CPUPercent但表达式 tag 是cpu——写指标表达式时可用的是cpu、memory、memoryUsage、mounts而不是 Go 字段名。同时可以推断映射过程是刻意裁剪的网络/磁盘计数属于累计值不适合做阈值告警故未暴露给表达式环境。mounts里只包含空间统计成功Available true的挂载点无法测量的卷Windows 卷、权限错误会被跳过避免误触发或误抑制告警。DispatcherConfig 与 CloudConfig 的字段拷贝DispatcherConfig在内部与types包各有一份定义internal/notification/types.go 与 types/notification.go广播时需要逐字段拷贝。历史缺陷正是这里APIKey、Prefix、ExpiresAt这类 Cloud 字段在转换时被漏掉。当前的修复方式是职责分离见 multi_host_service.go// broadcastNotificationConfig 中 // Cloud dispatchers are excluded; cloud config is broadcast separately. for _, d : range notifDispatchers { if d.Type cloud { continue } ... }通知配置广播显式跳过 cloud 类型Cloud 的敏感字段改由broadcastCloudConfig走独立通道完整拷贝cc types.CloudConfig{ APIKey: ncc.APIKey, Prefix: ncc.Prefix, ExpiresAt: ncc.ExpiresAt, StreamLogs: ncc.StreamLogs, }这种宁可分两条广播、也不在一个结构体里塞可选敏感字段的做法是该类漏字段缺陷的可推广解法。并发模型与原子状态记忆文档对并发模型的总结与当前源码完全吻合可以归纳为四条xsync.Map 贯穿全局subscriptions、dispatchersManager、activeStreams日志监听器都是xsync.Map订阅增删不会锁死整个处理循环订阅的运行时统计用原子类型TriggerCount是atomic.Int64LastTriggeredAt是atomic.Pointer[time.Time]。Manager 的UpdateSubscription在克隆订阅时要手动搬运原子值因为原子类型不能随结构体字面量复制// Preserve runtime stats (atomics cant be copied in struct literal) updated.TriggerCount.Store(sub.TriggerCount.Load()) updated.LastTriggeredAt.Store(sub.LastTriggeredAt.Load())这是UpdateSubscription采用克隆 Compute模式的原因Compute保证同一订阅的更新原子生效克隆保证其他协程读到的一直是完整快照。每容器冷却用 xsync.Map 记录MetricCooldowns/EventCooldowns都是containerID - time.Time冷却按容器隔离——同一个指标告警规则可以独立通知 N 个不同容器而同一容器在冷却期内只通知一次sendSem 限制并发发送容量 5 的加权信号量发送侧先获取许可再调 Dispatcher。仍在源码中的竞态TriggeredContainerIDs 懒初始化AddTriggeredContainertypes.go目前仍是检查-再创建的懒初始化func (s *Subscription) AddTriggeredContainer(id string) { if s.TriggeredContainerIDs nil { s.TriggeredContainerIDs xsync.NewMap[string, struct{}]() } s.TriggeredContainerIDs.Store(id, struct{}{}) }两个协程同时进入nil分支时可能各自创建一个 map先创建的那个被覆盖丢失其中的容器 ID。该路径由多条事件处理协程并发调用日志、指标、事件处理协程都可能触发告警统计是记忆文档标记的、且在当前 HEAD 上仍然存在的真实竞态风险点。config.go在从磁盘加载配置时会克隆TriggeredContainerIDsAddSubscription会预建MetricCooldowns等三张 map却唯独没有预建TriggeredContainerIDs——若希望消除竞态最直接的改进就是在AddSubscription/ReplaceSubscription中一并初始化它。前端校验日志表达式为什么是必填项记忆文档曾指出LogAlertFields的canSave允许空logExpression会创建死订阅。当前源码LogAlertFields.vue已收紧为必填注释解释了原因// An empty expression is not match everything — the backend treats a rule without a log // expression as inert — so it is a required field, not an optional filter. const canSave computed(() !!logExpression.value.trim() !logError.value);这与后端语义一致IsLogAlert()要求LogExpression ! LogProgram ! nil空表达式的日志订阅在后端是不产生任何匹配的惰性规则。前后端对空表达式的语义保持同构是避免死订阅的校验边界。多主机场景MultiHostService 的持久化与广播在单主机部署中Manager 之上还包了一层MultiHostService见 multi_host_service.go它持有ClientManager本地 Docker、远端 agent 客户端集合、通知 Manager 与磁盘Persister是 Web 层操作告警配置的统一入口func (m *MultiHostService) AddSubscription(sub *notification.Subscription) error { if err : m.notificationManager.AddSubscription(sub); err ! nil { return err } m.saveNotificationConfig() // SaveNotifications() broadcastNotificationConfig() return nil }StartNotificationManager的初始化顺序有讲究用本地 client构造三个监听器并创建 Manager记忆文档所述MultiHostService wraps Manager and handles config persistence agent broadcast即此先执行migration.MigrateCloudConfig把旧格式配置拆分为独立的 cloud.yml先Start()再Load()——注释写明Start first so matcher is available for LoadConfig因为加载订阅时ShouldListenToContainer匹配器必须可用广播一次已加载配置给所有已连接 agent并启动一个订阅主机上线事件的循环每当有新 host 变为Available重新广播通知与 Cloud 配置保证后加入的 agent 拿到当前完整规则集。Swarm 模式下还有swarmNotificationHandler同集群其他副本通过 agent 协议推送配置变更时它把变更写入本地 persister 并触发云客户端重连保证各副本的磁盘 内存锁步一致。通知配置的持久化位置是/data目录用户文档 docs/guide/alerts-and-webhooks.md 特别强调该目录必须挂载为卷否则重启后告警规则丢失。容器镜像更新container-update记录的风险点记忆文档还记录了一组针对容器镜像更新功能feat/container-update 特性分支的风险观察。需要说明这些条目属于特性分支时期的记录未能在当前仓库 HEAD 的通用路径中全部逐一对应验证此处仅作为该功能后续演进的风险清单转述progressCh 关闭契约不一致Docker 与 Agent 实现通过defer关闭progressCh而 K8s 实现没有——消费侧for range progressCh会永久挂起。跨实现共享的 channel 关闭契约需要每个实现显式履行NetworkSettings空指针风险docker.InspectResponse.NetworkSettings是指针ContainerCreate访问.Networks前缺少 nil 检查破坏性重建无回滚stop → remove → create → start 序列中若 create 在 remove 之后失败容器停留在已删除未重建状态没有回滚路径前端 SSE 解析使用手工 ReadableStream reader 而非 EventSource组件卸载时缺少 AbortController 清理可能造成流泄漏。这几条的共同模式是多实现/多阶段流程中的隐含契约与前述类型映射、channel 关闭检查属于同一类系统性风险审查此类代码时应优先核对契约的每个履行点。小结审查该子系统的核对清单结合记忆文档与源码核对的结果维护 Dozzle 通知系统时可按以下清单逐条验证核对项源码位置当前状态订阅/投递器容器均为并发安全 mapmanager.go成立xsync.Map指标/事件监听器随订阅启停updateListenersmanager.go成立enrich检查 stats channel 关闭stats_listener.go已含if !ok { return }Cloud 敏感字段独立广播、不漏字段multi_host_service.go已分离为broadcastCloudConfig表达式求值环境字段与表达式 tag 一致types/notification.go表达式用cpu/memory等 tag前端日志表达式必填LogAlertFields.vue已收紧为必填TriggeredContainerIDs懒初始化竞态types.go仍存在于当前 HEAD原子字段克隆时手动搬运manager.go成立克隆订阅时不可遗漏这套文档记录风险 → 源码逐条核对 → 以当前实现为准修正结论的工作方式本身就展示了 Dozzle 仓库如何用 Agent Memory 文件沉淀架构认知记忆文档是线索而非结论每一条都应回到 internal/notification 与 internal/support/docker 的源码中验证再决定它是待修缺陷还是历史痕迹。【免费下载链接】dozzleRealtime log viewer for containers. Supports Docker, Swarm and K8s.项目地址: https://gitcode.com/GitHub_Trending/do/dozzle创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考