Alex Rivera | Logout

F# async stack overflow

Asked 2011-08-06T21:58:53.537
9

I am surprised by a stack overflow in my async-based program. I suspect the main problem is with the following function, which is supposed to compose two async computations to execute in parallel and wait for both to finish:

let ( <|> ) (a: Async<unit>) (b: Async<unit>) =
    async {
        let! x = Async.StartChild a
        let! y = Async.StartChild b
        do! x
        do! y
    }

With this defined, I have the following mapReduce program that attempts to exploit parallelism in both the map and the reduce part. Informally, the idea is to spark N mappers and N-1 reducers using a shared channel, wait for them to finish, and read the result from the channel. I had my own Channel implementation, here replaced by a ConcurrentBag for shorter code (the problem affects both):

let mapReduce (map    : 'T1 -> Async<'T2>)
              (reduce : 'T2 -> 'T2 -> Async<'T2>)
              (input  : seq<'T1>) : Async<'T2> =
    let bag = System.Collections.Concurrent.ConcurrentBag()

    let rec read () =
        async {
            match bag.TryTake() with
            | true, value -> return value
            | _           -> do! Async.Sleep 100
                             return! read ()
        }

    let write x =
        bag.Add x
        async.Return ()

    let reducer =
        async {
            let! x = read ()
            let! y = read ()
            let! r = reduce x y
            return bag.Add r
        }

    let work =
        input
        |> Seq.map (fun x -> async.Bind(map x, write))
        |> Seq.reduce (fun m1 m2 -> m1 <|> m2 <|> reducer)

    async {
        do! work
        return! read ()
    }

Now the following basic test starts to throw StackOverflowException on n=10000:

let test n  =
    let m
Edit
Report

1 Answer

2

Very interesting discussion! I had a similar issue with Async.Parallel

let (<||>) first second = async { let! results = Async.Parallel([|first; second|]) in return   (results.[0], results.[1]) } 

let test = async { do! Async.Sleep 100 } 
(test, [1..10000]) 
||> List.fold (fun state value -> (test <||> state) |> Async.Ignore) 
|> Async.RunSynchronously // stackoverflow

I was very frustrated... so I solved it by creating my own Parallel combinator.

let parallel<'T>(computations : Async<'T> []) : Async<'T []> =
  Async.FromContinuations (fun (cont, exnCont, _) ->
    let count = ref computations.Length
    let results : 'T [] = Array.zeroCreate computations.Length
    computations 
        |> Array.iteri (fun i computation ->
            Async.Start <|
                async { 
                    try
                        let! res = computation
                        results.[i] <- res 
                    with ex -> exnCont ex

                    let n = System.Threading.Interlocked.Decrement(count)
                    if n = 0 then 
                        results |> cont 
                }))

And finally inspired by the discussion, I implemented the following mapReduce function

// (|f ,⊗|)

let mapReduce (mapF : 'T -> Async<'R>) (reduceF : 'R -> 'R -> Async<'R>) (input : 'T []) : Async<'R> = 
let rec mapReduce' s e =
    async { 
        if s + 1 >= e then return! mapF input.[s]
        else 
            let m = (s + e) / 2
            let! (left, right) =  mapReduce' s m <||> mapReduce' m e
            return! reduceF left right
    }
mapReduce' 0 input.Length
answered 2011-08-08T11:37:43.980

Your Answer