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 }