【问题标题】:Work queues in ClojureClojure 中的工作队列
【发布时间】:2012-09-12 01:02:22
【问题描述】:

我正在使用 Clojure 应用程序从 Web API 访问数据。我将发出大量请求,其中许多请求会导致发出更多请求,因此我希望将请求 URL 保留在队列中,以便在后续下载之间留出 60 秒。

按照this blog post我把这个放在一起:

(def queue-delay (* 1000 60)) ; one minute

(defn offer!
  [q x]
  (.offerLast q x)
  q)

(defn take!
  [q]
  (.takeFirst q))

(def my-queue (java.util.concurrent.LinkedBlockingDeque.))

(defn- process-queue-item
  [item]
  (println ">> " item)   ; this would be replaced by downloading `item`
  (Thread/sleep queue-delay))

如果我在代码中的某处包含(future (process-queue-item (take! my-queue))),那么在REPL 中我可以(offer! my-queue "something") 并立即看到“>> something”打印出来。到目前为止,一切都很好!但是我需要队列在我的程序处于活动状态的整个过程中持续存在。我刚才提到的(future ...) 调用可以将一个项目从队列中拉出,一旦它可用,但我想要一些能够持续观察队列并在有可用时调用process-queue-item

另外,与 Clojure 对并发的通常喜爱相反,我想确保一次只发出一个请求,并且我的程序等待 60 秒以发出每个后续请求。

我认为this Stack Overflow question 是相关的,但我不确定如何调整它来做我想做的事。如何连续轮询我的队列并确保一次只运行一个请求?

【问题讨论】:

  • 为什么要连续轮询,但每 60 秒才发送一次?每 60 秒轮询一次会完成同样的事情吗?
  • @mamboking 差不多,是的。这种方法的唯一缺点是将第一个项目添加到队列中:如果程序需要 5 秒钟来确定第一个请求 URL 将是什么,那么它将在那里等待 55 秒,直到检查队列。无论如何,该程序将运行很长时间,所以我想这不是什么大问题。
  • 您是否避免使用任务调度程序?例如,github.com/zcaudate/cronj(在该 repo 的自述文件中还有其他库的列表)
  • @georgek 我不一定要避免这样的事情,尽管这对我的应用程序来说似乎有点矫枉过正。

标签: clojure queue


【解决方案1】:

这是来自a project I did for fun 的代码 sn-p。它并不完美,但可以让您了解我是如何解决“等待 55 秒等待第一项”问题的。它基本上循环通过承诺,使用期货立即处理事情或直到承诺“变得”可用。

(defn ^:private process
  [queues]
  (loop [[q & qs :as q+qs] queues p (atom true)]
    (when-not (Thread/interrupted)
      (if (or
            (< (count (:promises @work-manager)) (:max-workers @work-manager))
            @p) ; blocks until a worker is available
        (if-let [job (dequeue q)]
          (let [f (future-call #(process-job job))]
            (recur queues (request-promise-from-work-manager)))
          (do
            (Thread/sleep 5000)
            (recur (if (nil? qs) queues qs) p)))
        (recur q+qs (request-promise-from-work-manager))))))

也许你可以做类似的事情?代码不是很好,可能需要重新编写以使用lazy-seq,但这只是我尚未完成的练习!

【讨论】:

    【解决方案2】:

    这很可能很疯狂,但你总是可以使用这样的函数来创建一个慢下来的惰性序列:

    (defn slow-seq [delay-ms coll]
      "Creates a lazy sequence with delays between each element"
      (lazy-seq 
        (if-let [s (seq coll)]
            (do 
              (Thread/sleep delay-ms)
              (cons (first s)
                    (slow-seq delay-ms (rest s)))))))
    

    这将基本上确保每个函数调用之间的延迟。

    您可以将它与以下内容一起使用,以毫秒为单位提供延迟:

    (doseq [i (slow-seq 500 (range 10))]
      (println (rand-int 10))
    

    或者,您也可以将函数调用放入序列中,例如:

    (take 10 (slow-seq 500 (repeatedly #(rand-int 10))))
    

    显然,在上述两种情况下,您都可以将(rand-int 10) 替换为您用于执行/触发下载的任何代码。

    【讨论】:

    • 如果我没看错的话,coll 的所有元素都必须在运行slow-seq 之前知道,对吧?我想要一些可以让你毫无问题地动态添加项目的东西。具体来说,如果一个 API 调用的结果是我需要进行另一个 API 调用,该函数是否允许将第二个调用放在队列中?
    【解决方案3】:

    我最终推出了自己的小型图书馆,我称之为simple-queue。你可以在 GitHub 上阅读完整的文档,但这里是完整的源代码。 我不会更新这个答案,所以如果你想使用这个库,请从 GitHub 获取源代码。

    (ns com.github.bdesham.simple-queue)
    
    (defn new-queue
      "Creates a new queue. Each trigger from the timer will cause the function f
      to be invoked with the next item from the queue. The queue begins processing
      immediately, which in practice means that the first item to be added to the
      queue is processed immediately."
      [f & opts]
      (let [options (into {:delaytime 1}
                          (select-keys (apply hash-map opts) [:delaytime])),
            delaytime (:delaytime options),
            queue {:queue (java.util.concurrent.LinkedBlockingDeque.)},
            task (proxy [java.util.TimerTask] []
                   (run []
                     (let [item (.takeFirst (:queue queue)),
                           value (:value item),
                           prom (:promise item)]
                       (if prom
                         (deliver prom (f value))
                         (f value))))),
            timer (java.util.Timer.)]
        (.schedule timer task 0 (int (* 1000 delaytime)))
        (assoc queue :timer timer)))
    
    (defn cancel
      "Permanently stops execution of the queue. If a task is already executing
      then it proceeds unharmed."
      [queue]
      (.cancel (:timer queue)))
    
    (defn process
      "Adds an item to the queue, blocking until it has been processed. Returns
      (f item)."
      [queue item]
      (let [prom (promise)]
        (.offerLast (:queue queue)
                    {:value item,
                     :promise prom})
        @prom))
    
    (defn add
      "Adds an item to the queue and returns immediately. The value of (f item) is
      discarded, so presumably f has side effects if you're using this."
      [queue item]
      (.offerLast (:queue queue)
                  {:value item,
                   :promise nil}))
    

    使用此队列返回值的示例:

    (def url-queue (q/new-queue slurp :delaytime 30))
    (def github (q/process url-queue "https://github.com"))
    (def google (q/process url-queue "http://www.google.com"))
    

    q/process 的调用将被阻塞,因此两个def 语句之间会有30 秒的延迟。

    使用此队列纯粹用于副作用的示例:

    (defn cache-url
      [{url :url, filename :filename}]
      (spit (java.io.File. filename)
            (slurp url)))
    
    (def url-queue (q/new-queue cache-url :delaytime 30))
    (q/add url-queue {:url "https://github.com",
                      :filename "github.html"})    ; returns immediately
    (q/add url-queue {:url "https://google.com",
                      :filename "google.html"})    ; returns immediately
    

    现在对q/add 的调用会立即返回。

    【讨论】:

      猜你喜欢
      • 2013-01-18
      • 2018-06-20
      • 2011-03-09
      • 2018-09-11
      • 2013-05-30
      • 2016-04-22
      • 1970-01-01
      • 2011-02-05
      • 2018-07-30
      相关资源
      最近更新 更多