use threading macro to tidy consumer.zk/messages
This commit is contained in:
parent
f9574228db
commit
1fca1666ec
1 changed files with 1 additions and 1 deletions
|
@ -39,7 +39,7 @@
|
|||
(let [[queue-seq queue-put] (pipe)]
|
||||
(doseq [[topic streams] (.createMessageStreams consumer (topic-map topics))]
|
||||
(future (doseq [msg (iterator-seq (.iterator (first streams)))]
|
||||
(queue-put (assoc (to-clojure msg) :topic topic)))))
|
||||
(queue-put (-> msg to-clojure (assoc :topic topic))))))
|
||||
queue-seq))
|
||||
|
||||
(defn topics
|
||||
|
|
Loading…
Reference in a new issue