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+
- ->timeunit
- deque
- deque?
- drain
- peek
- peek-first
- peek-last
- pop
- process-bulk
- push
- put
- put-first
- put-last
- queue
- queue:fixed
- queue:limited
- queue?
- remaining-capacity
- remove
- take
- take-first
- take-last
v 3.0
(defn ->timeunit
(^TimeUnit [obj]
(cond (keyword? obj)
(get +units+ obj)
:else
obj)))
link
(q/->timeunit :ms) => java.util.concurrent.TimeUnit/MILLISECONDS
v 3.0
(defn ^LinkedBlockingDeque deque
([]
(LinkedBlockingDeque.))
([& elements]
(LinkedBlockingDeque. ^java.util.Collection elements)))
link
(q/deque 1 2 3) => java.util.concurrent.LinkedBlockingDeque
v 3.0
(defn deque?
([obj]
(instance? BlockingDeque obj)))
link
(q/deque? (q/deque)) => true (q/deque? (q/queue)) => false
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]]
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]]
v 3.0
(defn peek-first
([^BlockingDeque queue]
(.peekFirst queue)))
link
(q/peek-first (q/deque 1 2 3)) => 1
v 3.0
(defn peek-last
([^BlockingDeque queue]
(.peekLast queue)))
link
(q/peek-last (q/deque 1 2 3)) => 3
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]]
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]
v 3.0
(defn ^BlockingQueue queue
([]
(LinkedBlockingQueue.))
([& elements]
(LinkedBlockingQueue. ^java.util.Collection elements)))
link
(q/queue 1 2 3) => java.util.concurrent.LinkedBlockingQueue
v 3.0
(defn ^BlockingQueue queue:fixed
([size]
(ArrayBlockingQueue. size)))
link
(type (q/queue:fixed 10)) => java.util.concurrent.ArrayBlockingQueue
v 3.0
(defn ^LimitedQueue queue:limited
([size]
(LimitedQueue. size)))
link
(q/queue? (q/queue:limited 10)) => true
v 3.0
(defn queue?
([obj]
(instance? BlockingQueue obj)))
link
(q/queue? (q/queue)) => true (q/queue? []) => false
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
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