排查有状态转换器join-averages返回结果异常问题
问题分析与修复
你的join-averages有状态转换器无法输出完整关联结果的核心问题在于没有正确遵循transducer的约定来传递最终的关联数据,反而在不该输出的时候错误推送了原始的单个输入项。
具体问题点
看这段关键逻辑:
(if (has-all-lists? c') c' ; 这里直接返回关联后的对象,但没通过reducing function(rf)更新累加器 (rf accumulator input)) ; 这里错误地把单个输入项推送到输出,导致日志里全是零散record
当某个uuid的三类数据收集齐全时,你直接返回了完整的关联对象,但transducer的规则是必须通过调用rf来更新accumulator,否则这个值根本不会进入输出通道。而在数据未收集齐时,调用(rf accumulator input)会把单个输入项推出去,这就是你看到大量零散record日志的原因。
修复后的完整代码
(:require [clojure.core.async :refer [chan sliding-buffer <! go-loop pipeline onto-chan] :as async] [clojure.set :refer [subset?]] [clojure.tools.logging :as log]) (defn has-all-lists? [averages-map] (subset? #{:tick-list :sma-list :ema-list} (->> averages-map keys (into #{})))) (defn join-averages [] (let [state (atom {})] (fn [rf] (fn ([] (rf)) ([accumulator] (rf accumulator)) ([accumulator input] (let [uuid (:uuid input) entry (cond (:last-trade-price-exponential input) {:ema-list input} (:last-trade-price-average input) {:sma-list input} (:last-trade-price input) {:tick-list input})] ;; 更新状态,合并当前输入项 (swap! state update uuid merge entry) (let [current (get @state uuid)] (log/info "c'" current) (if (has-all-lists? current) ;; 收集齐全时,通过rf把完整关联对象推入累加器(输出通道) (do ;; 可选:移除已完成的条目,避免内存泄漏 (swap! state dissoc uuid) (rf accumulator current)) ;; 未收集齐时,直接返回原累加器,不输出任何内容 accumulator)))))))) (comment (let [ema-list [{:uuid "1" :last-trade-price-exponential 10} {:uuid "2" :last-trade-price-exponential 11} {:uuid "3" :last-trade-price-exponential 12}] sma-list [{:uuid "1" :last-trade-price-average 10.1} {:uuid "2" :last-trade-price-average 10.2} {:uuid "3" :last-trade-price-average 10.3}] tick-list [{:uuid "1" :last-trade-price 11.1} {:uuid "2" :last-trade-price 11.2} {:uuid "3" :last-trade-price 11.3}] ec (chan (sliding-buffer 100)) sc (chan (sliding-buffer 100)) tc (chan (sliding-buffer 100)) _ (onto-chan ec ema-list) _ (onto-chan sc sma-list) _ (onto-chan tc tick-list) merged-ch (async/merge [tc sc ec]) output-ch (chan (sliding-buffer 100) (join-averages))] (async/pipeline 1 output-ch (join-averages) merged-ch) (go-loop [r (<! output-ch)] (when-not (nil? r) (log/info "record" r) (recur (<! output-ch))))))
关键修改说明
- 正确推送关联结果:当数据收集齐全时,通过
(rf accumulator current)把完整的关联对象传递给reducing function,这样它才会被推送到输出通道。 - 停止输出零散项:未收集齐数据时,直接返回原
accumulator,不调用rf推送单个输入项。 - 内存优化(可选):完成收集后从state中移除对应uuid的条目,避免长时间运行导致内存占用过高。
现在运行测试代码,你会看到输出通道打印出三个完整的关联对象,而不是零散的单个输入项了。
内容的提问来源于stack exchange,提问作者Nutritioustim
相关产品推荐
相关产品推荐

