std.concurrent.queue

blocking queues, deques, timeouts, and bulk draining

Use JVM blocking collections through a small Clojure API with consistent timeout units and queue-processing helpers.

1    Overview

The namespace constructs unbounded, fixed, limited, and double-ended blocking queues. Operations accept keyword time units such as :ms, :s, and :m.

2    Walkthrough

2.1    Produce and consume

(require '[std.concurrent.queue :as queue])

(def jobs (queue/queue:fixed 10))

(queue/put jobs {:id 1})
(queue/put jobs {:id 2} 100 :ms)
(queue/peek jobs)
(queue/take jobs 1 :s)
(queue/remaining-capacity jobs)

2.2    Drain work in batches

(queue/drain (queue/queue 1 2 3 4) 2)

(queue/process-bulk
 (fn [batch]
   (println "processing" batch))
 (queue/queue 1 2 3 4 5)
 3)

2.3    Use both ends

A deque supports access through put-first, put-last, take-first, and take-last.

(def work (queue/deque :middle))
(queue/put-first work :urgent)
(queue/put-last work :normal)
[(queue/take-first work)
 (queue/take-last work)]

3    API



+units+ ^

NONE
(def ^:private +units+)
link

->timeunit ^

[obj]
Added 3.0

returns the timeunit

v 3.0
(defn ->timeunit
  (^TimeUnit [obj]
   (cond (keyword? obj)
         (get +units+ obj)

         :else
         obj)))
link
(q/->timeunit :ms) => java.util.concurrent.TimeUnit/MILLISECONDS

deque ^

[] [& elements]
Added 3.0

constructs a blocking deque

v 3.0
(defn ^LinkedBlockingDeque deque
  ([]
   (LinkedBlockingDeque.))
  ([& elements]
   (LinkedBlockingDeque. ^java.util.Collection elements)))
link
(q/deque 1 2 3) => java.util.concurrent.LinkedBlockingDeque

deque? ^

[obj]
Added 3.0

checks if object is a `BlockingDeque`

v 3.0
(defn deque?
  ([obj]
   (instance? BlockingDeque obj)))
link
(q/deque? (q/deque)) => true (q/deque? (q/queue)) => false

drain ^

[queue] [queue max] [queue max target]
Added 3.0

drains elements to another vector

v 3.0
(defn drain
  ([^BlockingQueue queue]
   (let [target (java.util.ArrayList. (count queue))]
     (.drainTo queue target)
     target))
  ([^BlockingQueue queue ^long max]
   (let [target (java.util.ArrayList. max)]
     (drain queue max target)
     target))
  ([^BlockingQueue queue ^long max ^java.util.Collection target]
   (.drainTo queue target max)
   target))
link
((juxt #(q/drain % 2) #(into [] %)) (q/queue 1 2 3 4)) => [[1 2] [3 4]]

peek ^

[queue]
Added 3.0

takes element at the front of the queue

v 3.0
(defn peek
  ([^BlockingQueue queue]
   (.peek queue)))
link
(let [queue (q/queue 1 2 3)] [(q/peek queue) (vec queue)]) => [1 [1 2 3]]

peek-first ^

[queue]
Added 3.0

peeks from the front of the queue

v 3.0
(defn peek-first
  ([^BlockingDeque queue]
   (.peekFirst queue)))
link
(q/peek-first (q/deque 1 2 3)) => 1

peek-last ^

[queue]
Added 3.0

peeks from the back of the queue

v 3.0
(defn peek-last
  ([^BlockingDeque queue]
   (.peekLast queue)))
link
(q/peek-last (q/deque 1 2 3)) => 3

pop ^

[queue]
Added 3.0

pops from the front of the queue

v 3.0
(defn pop
  ([^BlockingDeque queue]
   (.pop queue)))
link
(q/pop (q/deque 1 2 3)) => 1

process-bulk ^

[f queue maximum]
Added 3.0

processes elements in the queue

v 3.0
(defn process-bulk
  ([f queue maximum]
   (loop [total  (count queue)]
     (let [n      (if (> total maximum)
                    maximum
                    total)
           txs (drain queue n)
           _   (when-not (empty? txs)
                 (f txs))]
       (let [more (count queue)]
         (if-not (zero? more)
           (recur more)))))))
link
(def +state+ (atom [])) (q/process-bulk (fn [elems] (swap! +state+ conj elems)) (q/queue 1 2 3 4 5) 3) @+state+ => [[1 2 3] [4 5]]

push ^

[queue element]
Added 3.0

puts at the front of the queue

v 3.0
(defn push
  ([^BlockingDeque queue element]
   (.push queue element)))
link
(-> (doto (q/deque) (q/push 1)) (vec)) => [1]

put ^

[queue element] [queue element timeout] [queue element timeout timeunit]
Added 3.0

puts an element at the back

v 3.0
(defn put
  ([^BlockingQueue queue element]
   (put queue element nil))
  ([^BlockingQueue queue element timeout]
   (put queue element timeout :ms))
  ([^BlockingQueue queue element timeout timeunit]
   (if (nil? timeout)
     (.put queue element)
     (.offer queue element timeout (->timeunit timeunit)))))
link
(->> (doto (q/queue) (q/put 1)) (into [])) => [1]

put-first ^

[queue element] [queue element timeout] [queue element timeout timeunit]
Added 3.0

puts at the front of the queue

v 3.0
(defn put-first
  ([^BlockingDeque queue element]
   (put-first queue element nil))
  ([^BlockingDeque queue element timeout]
   (put-first queue element timeout :ms))
  ([^BlockingDeque queue element timeout timeunit]
   (if (nil? timeout)
     (.putFirst queue element)
     (.offerFirst queue element timeout (->timeunit timeunit)))))
link
(-> (doto (q/deque) (q/put-first 1)) (vec)) => [1]

put-last ^

[queue element] [queue element timeout] [queue element timeout timeunit]
Added 3.0

puts at the back of the queue

v 3.0
(defn put-last
  ([^BlockingDeque queue element]
   (put-last queue element nil))
  ([^BlockingDeque queue element timeout]
   (put-last queue element timeout :ms))
  ([^BlockingDeque queue element timeout timeunit]
   (if (nil? timeout)
     (.putLast queue element)
     (.offerLast queue element timeout (->timeunit timeunit)))))
link
(-> (doto (q/deque) (q/put-last 1)) (vec)) => [1]

queue ^

[] [& elements]
Added 3.0

constructs a blocking queue

v 3.0
(defn ^BlockingQueue queue
  ([]
   (LinkedBlockingQueue.))
  ([& elements]
   (LinkedBlockingQueue. ^java.util.Collection elements)))
link
(q/queue 1 2 3) => java.util.concurrent.LinkedBlockingQueue

queue:fixed ^

[size]
Added 3.0

constructs a fixed size blocking queue

v 3.0
(defn ^BlockingQueue queue:fixed
  ([size]
   (ArrayBlockingQueue. size)))
link
(type (q/queue:fixed 10)) => java.util.concurrent.ArrayBlockingQueue

queue:limited ^

[size]
Added 3.0

constructs a limited queue

v 3.0
(defn ^LimitedQueue queue:limited
  ([size]
   (LimitedQueue. size)))
link
(q/queue? (q/queue:limited 10)) => true

queue? ^

[obj]
Added 3.0

checks if object is a `BlockingQueue`

v 3.0
(defn queue?
  ([obj]
   (instance? BlockingQueue obj)))
link
(q/queue? (q/queue)) => true (q/queue? []) => false

remaining-capacity ^

[queue]
Added 3.0

returns the remaining capacity

v 3.0
(defn remaining-capacity
  ([^BlockingQueue queue]
   (.remainingCapacity queue)))
link
(q/remaining-capacity (q/queue)) => 2147483647 (q/remaining-capacity (doto (q/queue:fixed 2) (q/put 1))) => 1

remove ^

[queue obj]
Added 3.0

removes element from queue

v 3.0
(defn remove
  ([^BlockingQueue queue obj]
   (.remove queue obj)))
link
(let [q (q/deque 1 2 3)] (q/remove q 2) (into [] q)) => [1 3]

take ^

[queue] [queue timeout] [queue timeout timeunit]
Added 3.0

takes an element from the queue

v 3.0
(defn take
  ([^BlockingQueue queue]
   (take queue nil))
  ([^BlockingQueue queue timeout]
   (take queue timeout :ms))
  ([^BlockingQueue queue timeout timeunit]
   (if (nil? timeout)
     (.take queue)
     (.poll queue timeout (->timeunit timeunit)))))
link
((juxt q/take #(into [] %)) (q/queue 1 2 3)) => [1 [2 3]]

take-first ^

[queue] [queue timeout] [queue timeout timeunit]
Added 3.0

takes from the front of the queue

v 3.0
(defn take-first
  ([^BlockingDeque queue]
   (.takeFirst queue))
  ([^BlockingDeque queue timeout]
   (take-first queue timeout :ms))
  ([^BlockingDeque queue timeout timeunit]
   (if (nil? timeout)
     (.takeFirst queue)
     (.pollFirst queue timeout (->timeunit timeunit)))))
link
(q/take-first (q/deque 1 2 3)) => 1

take-last ^

[queue] [queue timeout] [queue timeout timeunit]
Added 3.0

takes from the back of the queue

v 3.0
(defn take-last
  ([^BlockingDeque queue]
   (.takeLast queue))
  ([^BlockingDeque queue timeout]
   (take-last queue timeout :ms))
  ([^BlockingDeque queue timeout timeunit]
   (if (nil? timeout)
     (.takeLast queue)
     (.pollLast queue timeout (->timeunit timeunit)))))
link
(q/take-last (q/deque 1 2 3)) => 3