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

