missionary 2025-06-04

Hello folks! I have a question about parallel processing of multiple sequential flows. Looks like some big overhead arise, or me just didn't noticed something.

(ns example.core)

(comment
  (require '[missionary.core :as m])
  (def ls (take 300 (map #(array-map :id (rand-int 3) :sleep (rand-int 100) :number %) (range))))
                    ;;=> ({:id 0, :sleep 86, :number 0}
                    ;;    {:id 1, :sleep 63, :number 1}
                    ;;    ...)

  (reduce + (map :sleep ls))
  ;;=> 15334

  (defn process
    [msg]
    (m/via m/cpu
           (m/? (m/sleep (:sleep msg)))
           #_(println "processed" msg)
           msg))
  ;;=> #'exapmle.core/process

  (def run
    (let [flow (m/group-by :id (m/seed ls))]
      (m/ap
       (let [[k >msgs] (m/?> ##Inf flow)
             >msg (m/?> >msgs)]
         [k (m/? (process >msg))]))))
  ;;=> #'example.core/run

  ;; Execute
  (def result (time (m/? (m/reduce (fn [acc [k v]]
                                     (update acc k #(vec (conj % v)))) {} run))))
  ;;=> #'example.core/result

  result
  ;;=> {0 [{:id 0, :sleep 50, :number 0} {:id 0, :sleep 87, :number 1} ...],
  ;;    2 [{:id 2, :sleep 71, :number 2} {:id 2, :sleep 59, :number 4} ...],
  ;;    1 [{:id 1, :sleep 25, :number 5} {:id 1, :sleep 37, :number 11} ...]}

  (into {} (map (fn [[k v]] [k (reduce + (map :sleep v))]) result))
  ;;=> {0 4630, 2 4899, 1 5805}
  )
Execution of (def result (time (m/? (m/reduce ... took "Elapsed time: 9705.874292 msecs" As we can see, total amount of sleep is 15334 ms. Something like that it must take if we process all ls items sequentially. But processing is separated for items with different :id , in this case three different threads. Like it is printed in last expression, sums of sleeps for every thread is not greated then 5805. Yes, I do expect some overhead of forking etc... but not in 9705 / 5805 ~= 1.67 times This coefficient differ with different amount of ls items and rand-int arguments, but always much greater the one, actually 1.67 is even a good value... What is the problem?

Try to buffer each processing group, the execution time should be much closer to what you expect.

(def run
  (let [flow (m/group-by :id (m/seed ls))]
    (m/ap
      (let [[k >msgs] (m/?> ##Inf flow)
            >msg (m/?> (m/buffer 10 >msgs))]
        [k (m/? (process >msg))]))))
What you observe here is the backpressure of group-by - soon after boot the 3 group consumers are saturated and start to propagate backpressure which pauses consumption of input seed. Now if the next input item is e.g. id 0, group-by will wait for consumer group 0 to become ready to process, even if consumer groups 1 and 2 become available before. The problem will likely get worse if ids are not randomly distributed.

👍 1

Thank you for explanation! That definitely makes sense and works. Is there any mechanism to handle buffer filling? Like to run side-effect (e.g. logging or alert) if capacity is reached...

There's no such mechanism available. If you want to roll your own, the starting point would be to add a side-effecting stage before each consumer group to detect a value transfer, then compare timings somehow.

👌 1

Okay, I will take a look

How come this flow here only ever uses 0 or 1 threads? I have an e/watch on !store-threads and I see that it's usually 0, sometimes 1.

(defn mk-store-documents-f [kview documents]
  (m/ap (let [doc (m/?> 3 documents)
              _ (swap! !store-threads inc)
              doc (m/? (mk-store-document-t kview doc))
              _ (swap! !store-threads dec)]
          doc)))
This stage in the pipeline is backpressuring the previous stages, which is fine, but when the rest of the system is halted I would think this stage should have as many threads as it needs?

I'm just exploring this lib after using electric for a bit and for now I use m/via just like I used to use clojure.core.async/thread if that helps

In my Electric UI I render the number of store-documents threads and the size of the buffer before store-documents-f. 0 store threads, 13 in the buffer.

Here's the pipeline as a whole

(m/?
          (m/reduce
           rfs/last
           (->>
            kudos/documents
            (mk-filter-documents-f kview)
            #_(m/buffer 512)
            (mk-prepare-documents-f kview)
            (m/eduction (map #(do (swap! !buffer-before-store-documents-size inc) %)))
            (m/buffer 128)
            (m/eduction (map #(do (swap! !buffer-before-store-documents-size dec) %)))
            (mk-store-documents-f kview))))

have you checked for blocking calls in mk-store-document ?

There's a blocking call within mk-store-document, should I wrap that in a m/via?

Ah, right I think I get it maybe. Thank you!

I think Missionary clicked for me now 😎 So (m/via m/blk <form>) is what we use to glue blocking-IO code into a Missionary tasks. And Missionary tasks are always async. m/sp does not do anything to prevent blocking forms from making the task as a whole from blocking. Tasks must never block. The m/sp macro gives us the ability to park on tasks instead of blocking. Parking and blocking has the same semantics, but different runtime characteristics. Blocking consumes a thread while waiting, whereas parking does not 🛣️ m/sp also gives us cancellation semantics. When an m/sp task t1 is executed, it executes synchronously until a parking form (m/? t2) and then it waits asynchronously until either t2 completes or t1 is cancelled. If t1 is cancelled, t2 is also automatically cancelled. WRONG:

(m/sp (http-get ""))
RIGHT:
(m/sp (m/? (m/via m/blk (http-get ""))))
NICER:
(defn mk-http-get-t [req]
  (m/via m/blk (http-get req)))

(m/sp (m/? (mk-http-get-t "")))
Is that right? Sorry for being slow, this async stuff is super hard for me 😅

📝 2
👌 2

That is 100% correct 👍

Fantastic! 😄