Empresas
Empregos
  • Sobre nós
  • Soluções
    • Publicação de vagas
      Publique sua vaga e receba candidatos qualificados em 48h.
    • Avaliações de candidatos
      Mais de 500 testes técnicos e psicológicos, mais anti-fraude.
    • Headhunting
      Busca executiva personalizada do início ao fim.
    • Folha de Pagamento + EOR
      Dispersão de folha e EOR em mais de 15 países da LATAM.
  • Preços
  • Empregos

0

106
Visualizações
F# desordenado AsyncSeq.mapParallel with throttling

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:

  1. use 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.
  2. use 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.
  3. construya lo que necesito enumerando la fuente y emitiendo resultados a través 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 maxDegreeOfParallelism

Seguramente 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...).

over 4 years ago · Santiago Trujillo
2 Respostas
Responde à pergunta

0

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 }
over 4 years ago · Santiago Trujillo Relatório

0

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.

over 4 years ago · Santiago Trujillo Relatório
Responde à pergunta
Encontrar trabalhos remotos

Descubra a nova forma de encontrar um emprego!

melhores empregos
Principais categorias de trabalho
Empresas
Postar vaga Preços Comercial
Jurídico
Termos e Condições Política de privacidade
© 2026 PeakU Inc. All Rights Reserved.
Andres GPT
Recomende algumas ofertas para mim
Preciso de ajuda