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

如何用ClojureScript core/async限制服务并发请求数至5?

控制core.async中服务请求的并发数

要解决这个并发限制的问题,我们可以利用core.async的固定Worker协程池或者专用流式处理函数来实现,下面给你两种实用方案,按需选择就行:

方案一:固定Worker协程池(直观易理解)

思路很简单:启动5个独立的Worker协程,它们共享同一个输入通道读取请求参数,同一时间最多只有5个Worker在发起服务请求,完美符合你的并发限制要求。

修改后的代码如下:

(require '[clojure.core.async :as async :refer [go go-loop >! <! chan]])

(defn produce-inputs [in-chan inputs]
  (let [input-names-seq (map #(:name %) inputs)]
    ;; 写完所有参数后关闭输入通道,让Worker知道没有新任务了
    (go
      (doseq [input-name input-names-seq]
        (>! in-chan input-name))
      (async/close! in-chan))))

(defn consume [inputs]
  (let [in-chan (async/chan 10)  ;; 适当增大缓冲,避免生产者阻塞
        out-chan (async/chan 10)]

    ;; 启动5个Worker协程,专门处理服务请求
    (dotimes [_ 5]
      (go-loop []
        (when-let [input-name (<! in-chan)]
          ;; 发起异步服务请求,回调中把结果写入输出通道
          (retrieve-resource-from-service 
            input-name
            (fn [resp]
              (go
                (let [result (:result resp)]
                  (>! out-chan result)))))
          (recur))))

    ;; 生产者写入所有请求参数
    (produce-inputs in-chan inputs)

    ;; 读取结果并执行CPU密集型任务
    (go-loop []
      (when-let [result (<! out-chan)]
        (do-some-cpu-heavy-work result)
        (recur)))

    ;; 返回输出通道,方便后续扩展处理(可选)
    out-chan))

;; 入口函数
(defn run [inputs]
  (consume inputs))

关键修改点:

  • 关闭输入通道:生产者完成参数写入后关闭in-chan,Worker协程会在通道关闭后自动退出循环,避免协程泄漏。
  • 固定Worker数量:用dotimes启动5个Worker,每个Worker循环处理请求,确保同时活跃的服务请求不超过5个。
  • 调整通道缓冲:增大通道缓冲可以避免生产者因通道满而阻塞,适合参数数量较多的场景。

方案二:使用pipeline-async(更简洁的流式处理)

core.async提供的pipeline-async专门用于异步转换的流式处理,直接指定并行度就能控制并发数,代码更简洁:

(require '[clojure.core.async :as async :refer [go >! <! chan pipeline-async]])

(defn produce-inputs [in-chan inputs]
  (go
    (doseq [input-name (map #(:name %) inputs)]
      (>! in-chan input-name))
    (async/close! in-chan)))

(defn consume [inputs]
  (let [in-chan (async/chan 10)
        out-chan (async/chan 10)]

    ;; 用pipeline-async控制并行度为5,处理服务请求
    (pipeline-async 5 out-chan
      (fn [input-name callback]
        ;; 发起异步请求,完成后通过callback返回结果
        (retrieve-resource-from-service input-name
          (fn [resp]
            (callback (list (:result resp))))))
      in-chan)

    (produce-inputs in-chan inputs)

    ;; 处理CPU密集型任务(建议放到单独线程避免阻塞协程池)
    (go-loop []
      (when-let [result (<! out-chan)]
        (async/thread (do-some-cpu-heavy-work result))
        (recur)))

    out-chan))

(defn run [inputs]
  (consume inputs))

额外优化提示:

因为do-some-cpu-heavy-work是CPU密集型任务,直接在go块里执行会占用core.async的线程池资源,建议用async/thread把它放到独立线程中执行,避免影响IO协程的调度效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 04:10:20