Estoy usando F# y tengo un AsyncSeq<Async<'t>> . Cada elemento tardará una cantidad variable de tiempo en procesarse y realiza E/S con velocidad limitada.
Quiero ejecutar todas las operaciones en paralelo y luego pasarlas por la cadena como AsyncSeq<'t> para poder realizar más manipulaciones en ellas y, en última instancia, AsyncSeq.fold en un resultado final.
Las siguientes operaciones AsyncSeq casi satisfacen mis necesidades:
mapAsyncParallel : hace el paralelismo, pero no tiene restricciones (y no necesito que se conserve el orden)iterAsyncParallelThrottled : paralelo y tiene un grado máximo de paralelismo pero no me permite devolver resultados (y no necesito que se conserve el orden) Lo que realmente necesito es como un mapAsyncParallelThrottled . Pero, para ser más precisos, realmente la operación se titularía mapAsyncParallelThrottledUnordered .
Cosas que estoy considerando:
mapAsyncParallel pero use un Semaphore dentro de la función para restringir el paralelismo, lo que probablemente no sea óptimo en términos de concurrencia, y debido al almacenamiento en búfer de los resultados para reordenarlos.iterAsyncParallelThrottled y haga un plegado feo de los resultados en un acumulador a medida que llegan protegidos por un candado como este , pero no necesito el pedido, por lo que no será óptimo.AsyncSeqSrc como este . Probablemente tendría un conjunto de tareas Async.StartAsTask en vuelo y comenzaría más después de cada Task.WaitAny me da algo para AsyncSeqSrc.put hasta que alcance el maxDegreeOfParallelismSeguramente me estoy perdiendo una respuesta simple y hay una mejor manera.
De lo contrario, me encantaría que alguien revisara mi opción 3 en cualquier dirección.
Estoy abierto a usar AsyncSeq.toAsyncEnum y luego uso una forma IAsyncEnumerable de lograr el mismo resultado si existe, aunque idealmente sin entrar en TPL DataFlow o RX land si se puede evitar (he hecho una búsqueda exhaustiva de SO sin resultados...).
Si entiendo sus requisitos, entonces algo como esto funcionará. Combina efectivamente el iter desordenado con un canal para permitir un mapeo en su lugar.
let mapAsyncParallelBoundedUnordered boundedAmount (mapper: 't -> Async<_>) source = asyncSeq { let! ct = Async.CancellationToken let channel = Channel.CreateUnbounded() let! _ = async { do! source |> AsyncSeq.iterAsyncParallelThrottled boundedAmount (fun s -> async { let! orderChild = mapper s do! channel.Writer.WriteAsync(orderChild, ct) }) channel.Writer.Complete() } |> Async.StartChild for item in channel.Reader.ReadAllAsync(ct) |> AsyncSeq.ofAsyncEnum do let! toReturn = item yield toReturn }También con un poco de variación de lo anterior (por ejemplo, tareas secundarias) puede hacerlo ordenado y delimitado por paralelismo.
let mapAsyncParallelBounded boundedAmount mapper source = asyncSeq { let! ct = Async.CancellationToken let channel = Channel.CreateBounded(BoundedChannelOptions(boundedAmount)) let! _ = source |> AsyncSeq.iterAsync (fun s -> async { let! orderChild = mapper s |> Async.StartChild do! channel.Writer.WriteAsync(orderChild, ct) }) |> Async.StartChild let! ct = Async.CancellationToken for item in channel.Reader.ReadAllAsync(ct) |> AsyncSeq.ofAsyncEnum do let! toReturn = item yield toReturn }Aquí hay un banco de pruebas que usé para validar el excelente trabajo de @akara:
#r "nuget:FSharp.Control.AsyncSeq" open FSharp.Control module AsyncSeqEx = open System.Threading.Channels let mapAsyncParallelBoundedUnordered boundedAmount (mapper: 't -> Async<'u>) source = asyncSeq { let! ct = Async.CancellationToken let channel : Channel<'u> = Channel.CreateUnbounded() let handle req = async { let! res = mapper req do! let t = channel.Writer.WriteAsync(res, ct) in t.AsTask() |> Async.AwaitTask } let! _ = Async.StartChild <| async { do! source |> AsyncSeq.iterAsyncParallelThrottled boundedAmount handle channel.Writer.Complete() } yield! channel.Reader.ReadAllAsync(ct) |> AsyncSeq.ofAsyncEnum } También transfirí el mismo código para usar AsyncSeqSrc en lugar de canales, lo que parece funcionar también, con un rendimiento equivalente:
// AsyncSeqSrc-based reimpl of the above let mapAsyncParallelBoundedUnordered2 boundedAmount (mapper: 't -> Async<'u>) source = asyncSeq { let output = AsyncSeqSrc.create () let handle req = async { let! res = mapper req in AsyncSeqSrc.put res output } let! _ = Async.StartChild <| async { do! source |> AsyncSeq.iterAsyncParallelThrottled boundedAmount handle AsyncSeqSrc.close output } yield! AsyncSeqSrc.toAsyncSeq output } El siguiente impl, apoyándose en AsyncSeq.mapAsyncParallel parece lograr un rendimiento similar a ambos:
module Async = let parallelThrottled dop f = Async.Parallel(f, maxDegreeOfParallelism = dop) type Semaphore(max) = let inner = new System.Threading.SemaphoreSlim(max) member _.Await() = async { let! ct = Async.CancellationToken return! inner.WaitAsync ct |> Async.AwaitTask } member _.Release() = inner.Release() |> ignore let throttle degreeOfParallelism f = let s = Semaphore degreeOfParallelism fun x -> async { do! s.Await() try return! fx finally s.Release() } module AsyncSeq = open FSharp.Control // see https://stackoverflow.com/a/71065152/11635 let mapAsyncParallelThrottled degreeOfParallelism (f: 't -> Async<'u>) : AsyncSeq<'t> -> AsyncSeq<'u> = let throttle = Async.throttle degreeOfParallelism AsyncSeq.mapAsyncParallel (throttle f)banco de pruebas:
let dop = 10 let r = System.Random() let durations = Array.init 10000 (fun _ -> r.Next(10, 100)) let work = let sleep (x : int) = async { do! Async.Sleep x return x } AsyncSeq.ofSeq durations |> AsyncSeq.mapAsyncParallelThrottled dop sleep let start = System.Diagnostics.Stopwatch.StartNew() let results = work |> AsyncSeq.toArrayAsync |> Async.RunSynchronously let timeTaken = start.ElapsedMilliseconds let totalTimeTaken = Array.sum results let expectedWallTime = float totalTimeTaken / float dop let overhead = timeTaken - int64 expectedWallTime let inline stringf format (x : ^a) = (^a : (member ToString : string -> string) (x, format)) let inline sep x = stringf "N0" x printfn $"Gross {sep totalTimeTaken}ms Threads {dop} Wall {sep timeTaken}ms overhead {sep overhead}ms ordered: {durations = results}"Resultado:
Bruto 544 873 ms Subprocesos 10 Muro 55 659 ms sobrecarga 1172 ms ordenado: Verdadero
Por ahora, parece que para mi caso de uso no hay una gran victoria al admitir resultados desordenados en lugar de solo tener el argumento de función para el autogobierno mapAsyncParallel para lograr el efecto de limitación deseado.