代码之家  ›  专栏  ›  技术社区  ›  varnie

Clojurescript:使用核心/异步通道以块的形式处理请求

  •  2
  • varnie  · 技术社区  · 8 年前

    我有以下场景:

    我使用一些服务来检索一些数据,并将其传递给我的输入。

    有了一些输入参数,我需要对上述服务执行N个请求,收集输出,并为每个输出执行一些CPU繁重的任务。

    我正在尝试使用核心/异步通道来实现这一点。

    这是我的尝试(示意图),它有点工作,但它没有表现出我想要的。 如有任何改进建议,我将不胜感激。

    (defn produce-inputs
      [in-chan inputs]
      (let input-names-seq (map #(:name %) inputs)]
        (doseq [input-name input-names-seq]
          (async/go
            (async/>! in-chan input-name)))))
    
    (defn consume
      [inputs]
      (let [in-chan (async/chan 1)
            out-chan (async/chan 1)]
            (do
              (produce-inputs in-chan inputs)
              (async/go-loop []
                       (let [input-name (async/<! in-chan)]
                         (do
                             (retrieve-resource-from-service input-name 
                                                             ; response handler
                                                             (fn [resp]
                                                               (async/go
                                                                 (let [result (:result resp)]
                                                                   (async/>! out-chan result)))))
                             (when input-name
                               (recur)))))
    
         ; read from out-chan and do some heavy work for each entry
         (async/go-loop []
                       (let [result (async/<! out-chan)]
                             (do-some-cpu-heavy-work result))))))
    
    ; entry point
    (defn run
      [inputs]
      (consume inputs))
    

    有没有办法更新它,使每一时刻的服务请求不超过五个( retrieve-resource-from-service )活动?

    如果我的解释不清楚,请提问,我会更新它。

    1 回复  |  直到 8 年前
        1
  •  5
  •   Aleph Aleph    8 年前

    您可以创建另一个通道作为令牌桶来限制请求的速率。

    See this link 例如,使用令牌桶进行每秒速率限制。

    要限制同时请求的数量,可以执行以下操作:

    (defn consume [inputs]
      (let [in-chan (async/chan 1)
            out-chan (async/chan 1)
            bucket (async/chan 5)]
        ;; ...
        (dotimes [_ 5] (async/put! bucket :token))
        (async/go-loop []
          (let [input-name (async/<! in-chan)
                token (async/<! bucket)]
            (retrieve-resource-from-service
              input-name 
              ; response handler
              (fn [resp]
                (async/go
                  (let [result (:result resp)]
                    (async/>! out-chan result)
                    (async/>! bucket token)))))
            (when input-name
              (recur))))
        ;; ...
        ))
    

    新频道, bucket ,并将五个项目放入其中。在触发请求之前,我们从桶中取出一个令牌,并在请求完成后将其放回。如果在 水桶 频道,我们必须等待其中一个请求完成。

    注意:这只是代码的草图,您可能需要更正它。特别是,如果您的数据库中有任何错误处理程序 retrieve-resource-from-service 函数,您也应该在出现错误时将令牌放回,以避免最终死锁。