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

排查有状态转换器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))))))

关键修改说明

  1. 正确推送关联结果:当数据收集齐全时,通过(rf accumulator current)把完整的关联对象传递给reducing function,这样它才会被推送到输出通道。
  2. 停止输出零散项:未收集齐数据时,直接返回原accumulator,不调用rf推送单个输入项。
  3. 内存优化(可选):完成收集后从state中移除对应uuid的条目,避免长时间运行导致内存占用过高。

现在运行测试代码,你会看到输出通道打印出三个完整的关联对象,而不是零散的单个输入项了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:46:42