如何在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
相关产品推荐
相关产品推荐

