How to schedule a list of functions `n` seconds apart with Clojure Core Async

Viewed 202

In my Clojure project I'm trying to make a list of http calls to an API that has a rate limiter that only allows n calls per minute. I want each of the responses to be returned once all the http calls are finished for further processing. I am new to Clojure's Core Async, but thought it would be a good fit, but because I need to run each call n seconds apart I am also trying to use the Chime library. In Chime's library it has examples using Core Async, but the examples all call the same function at each time interval which won't work for this use case.

While there is probably a way to use chime-async that better serves this use case, all of my attempts at that have failed so I've tried simply wrapping Chime calls with core async, but I am probably more baffled by Core Async than Chime.

This is an example of my name space.

(ns mp.util.schedule
  (:require [chime.core :as chime]
            [clojure.core.async :as a]
            [tick.alpha.api :as tick]))

(defn schedule-fns
  "Takes a list of functions and a duration in seconds then runs each function in the list `sec` seconds apart
   optionally provide an inst to start from"
  [fs sec & [{:keys [inst] :or {inst (tick/now)}}]]
  (let [ch (a/chan (count fs))
        chime-times (map-indexed
                      (fn mapped-fn [i f]
                        (a/put! ch (chime/chime-at [(.plusSeconds inst (* i sec))]
                                                   (fn wrapped-fn [_] (f)))))
                      fs)]
    (doseq [chi chime-times]
      (a/<!! chi))))
; === Test Code ===

; simple test function
(defn sim-fn
  "simple function that prints a message and value, then returns the value"
  [v m]
  (println m :at (tick/now))
  v)

; list of test functions
(def fns [#(sim-fn 1 :one)
          #(sim-fn 2 :two)
          #(sim-fn 3 :three)])

What I want to happen when calling (schedule-fns fns 2) is for each function in fns to run n seconds from each other and for schedule-fns to return (1 2 3) (the return values of the functions), but this isn't what it is doing. It is calling each of the functions at the correct times (which I can see from the log statements) but it isn't returning anything and there's an error I don't understand. I'm getting:

(schedule-fns fns 2)
:one :at #time/instant "2021-03-05T23:31:52.565Z"
Execution error (IllegalArgumentException) at clojure.core.async.impl.protocols/eval11496$fn$G (protocols.clj:15).
No implementation of method: :take! of protocol: #'clojure.core.async.impl.protocols/ReadPort found for class: java.lang.Boolean
:two :at #time/instant "2021-03-05T23:31:54.568Z"
:three :at #time/instant "2021-03-05T23:31:56.569Z"

If I could get help getting my code to use Core Async properly (with or without Chime) I'd really appreciate it. Thanks.

4 Answers

Try this:

(defn sim-fn
  "simple function that prints a message and value, then returns the value"
  [v m]
  (println m)
  v)

; list of test functions
(def fns [#(sim-fn 1 :one)
          #(sim-fn 2 :two)
          #(sim-fn 3 :three)])

(defn schedule-fns [fns sec]
       (let [program (interpose #(Thread/sleep (* sec 1000))
                                fns)]
         (remove #(= % nil)
                 (for [p program]
                         (p)))))

Then call:

> (schedule-fns fns 2)
:one
:two
:three
=> (1 2 3)

I came up with a way to get what I want...with some caveats.

(def results (atom []))

(defn schedule-fns
  "Takes a list of functions and a duration in seconds then runs each function in the list `sec` seconds apart
   optionally provide an inst to start from"
  [fs sec]
  (let [ch (chan (count fs))]
    (go-loop []
      (swap! results conj (<! ch))
      (recur))
    (map-indexed (fn [i f]
                   (println :waiting (* i sec) :seconds)
                   (go (<! (timeout (* i sec 1000)))
                       (>! ch (f))))
                 fs)))

This code has the timing and behavior that I want, but I have to use an atom to store the responses. While I can add a watcher to determine when all the results are in, I still feel like I shouldn't have to do that.

I guess I'll use this for now, but at some point I'll keep working on this and if anyone has something better than this approach I'd love to see it.

I had a couple friends look at this and they each came up with different answers. These are certainly better than what I was doing.

(defn schedule-fns [fs secs]
  (let [ret (atom {})
        sink (a/chan)]
    (doseq [[n f] (map-indexed vector fs)]
      (a/thread (a/<!! (a/timeout (* 1000 n secs)))
                (let [val (f)
                      this-ret (swap! ret assoc n val)]
                  (when (= (count fs) (count this-ret))
                    (a/>!! sink (mapv (fn [i] (get this-ret i)) (range (count fs))))))))
    (a/<!! sink)))

and

(defn schedule-fns
  [fns sec]
  (let [concurrent (count fns)
        output-chan (a/chan)
        timedout-coll (map-indexed (fn [i f]
                                     #(do (println "Waiting")
                                          (a/<!! (a/timeout (* 1000 i sec)))
                                          (f))) fns)]
    (a/pipeline-blocking concurrent
                         output-chan
                         (map (fn [f] (f)))
                         (a/to-chan timedout-coll))
    (a/<!! (a/into [] output-chan))))

If your objective is to work around the rate limiter, you can consider implementing it in the async channel. Below is one sample implementation - the function takes a channel, throttled its input with a token based limiter and pipe it to an output channel.

(require '[clojure.core.async :as async])

(defn rate-limiting-ch [input xf rate]
  (let [tokens (numerator rate)
        period (denominator rate)
        ans    (async/chan tokens xf)
        next   (fn [] (+ period (System/currentTimeMillis)))]
    (async/go-loop [c tokens
                    t (next)]
      (if (zero? c)
        (do
          (async/<! (async/timeout (- t (System/currentTimeMillis))))
          (recur tokens (next)))
        (when-let [x (async/<! input)]
          (async/>! ans x)
          (recur (dec c) t))))
    ans))

And here is a sample usage:

(let [start  (System/currentTimeMillis)
      input  (async/to-chan (range 10))
      output (rate-limiting-ch input
                               ;; simulate an api call with roundtrip time of ~300ms
                               (map #(let [wait (rand-int 300)
                                           ans  {:time  (- (System/currentTimeMillis) start)
                                                 :wait  wait
                                                 :input %}]
                                       (Thread/sleep wait)
                                       ans))
                               ;; rate limited to 2 calls per 1000ms
                               2/1000)]
  ;; consume the output
  (async/go-loop []
    (when-let [x (async/<! output)]
      (println x)
      (recur))))

Output:

{:time 4, :wait 63, :input 0}
{:time 68, :wait 160, :input 1}
{:time 1003, :wait 74, :input 2}
{:time 1079, :wait 151, :input 3}
{:time 2003, :wait 165, :input 4}
{:time 2169, :wait 182, :input 5}
{:time 3003, :wait 5, :input 6}
{:time 3009, :wait 18, :input 7}
{:time 4007, :wait 138, :input 8}
{:time 4149, :wait 229, :input 9}
Related