You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Riemann Clojure中避免使用reinject并重构事件处理规则?

问题场景

我需要实现以下Riemann事件处理逻辑:

  • 对服务a的2个事件求和,生成a-derivative事件
  • 对服务b的3个事件求和,生成b-derivative事件
  • 计算这两个衍生事件的商

初始实现用了reinject重新注入衍生事件,再通过project计算商,但Riemann文档建议避免使用reinject。我尝试了用pipe+splitp的重构方案,但只适用于过滤条件类似的场景。如果事件处理逻辑完全不同(比如一个基于service名称,一个基于tags),该怎么处理?

初始实现(使用reinject)

(streams
 (where (service "a")
                #(info "a-" %)
                (moving-event-window 2 (smap folds/sum (with :service "a-derivative" reinject))))
  (where (service "b")
                #(info "b-" %)
                (moving-event-window 3 (smap folds/sum (with :service "b-derivative" reinject))))
 (project [(service "a-derivative") (service "b-derivative")]
          (smap folds/quotient (with :laa "doooo" prn))))

仅适用于相似过滤条件的重构方案

(streams
 (pipe -
       (splitp = service
               "a" (moving-event-window 2 (smap folds/sum (with :service "a-derivative" -)))
               "b" (moving-event-window 2 (smap folds/sum (with :service "b-derivative" -)))) 
            
       (sdo #(info "hmm" %)
            (project [(service "a-derivative") (service "b-derivative")]
                     (smap folds/quotient (with :laa "doooo" prn))))
       
       )
 )

针对不同过滤条件的解决方案

可以通过merge操作将多个独立处理分支的输出合并,再传递给project进行后续计算,完全不需要reinject。

比如假设一个分支基于service="a",另一个分支基于tags包含"b-tag",逻辑如下:

(streams
  (merge
    ;; 分支1:处理service="a"的事件,生成a-derivative
    (where (service "a")
      #(info "a-" %)
      (moving-event-window 2 (smap folds/sum (with :service "a-derivative" -))))
    ;; 分支2:处理tags包含"b-tag"的事件,生成b-derivative
    (where (tags #"b-tag")
      #(info "b-" %)
      (moving-event-window 3 (smap folds/sum (with :service "b-derivative" -))))
  )
  ;; 对合并后的衍生事件做project计算
  (project [(service "a-derivative") (service "b-derivative")]
    (smap folds/quotient (with :laa "doooo" prn))))

方案说明

  • merge会将多个流的输出合并成一个单一的流,不管每个分支的过滤条件和处理逻辑差异多大,只要最终输出的事件带有标记(比如这里的:service字段),就能被后续的project正确识别。
  • 每个分支可以完全独立定义过滤规则(如service匹配、tags匹配,甚至更复杂的自定义条件)和处理逻辑(不同的窗口大小、聚合函数等)。
  • 全程不需要reinject,符合Riemann的最佳实践,避免了事件重新注入带来的潜在问题(比如循环处理、事件重复等)。

内容的提问来源于stack exchange,提问作者Vipin Menon

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 06:15:59