Estoy tratando de ejecutar cierta función de suspensión varias veces, de tal manera que nunca se ejecuten más de N de estos al mismo tiempo.
Para aquellos familiarizados con las bibliotecas Akka y Scala Streaming, algo como mapAsync .
Hice mi propia implementación usando un canal de entrada (como en los canales de kotlin) y N canales de salida. Pero parece engorroso y no muy eficiente.
El código que estoy usando actualmente es algo así:
val inChannel = Channel<T>() val outChannels = (0..n).map{ Channel<T>() } launch{ var i = 0 for(t in inChannel){ outChannels[i].offer(t) i = ((i+1)%n) } } outChannels.forEach{outChannel -> launch{ for(t in outChannel){ fn(t) } } }Por supuesto que tiene gestión de errores y todo, pero aún así...
Editar: Hice la siguiente prueba y falló.
test("Parallelism is correctly capped") { val scope = CoroutineScope(Dispatchers.Default.limitedParallelism(3)) var num = 0 (1..100).map { scope.launch { num ++ println("started $it") delay(Long.MAX_VALUE) } } delay(500) assertEquals(3,num) }Puede usar la limitedParallelism en un Dispatcher (experimental en v1.6.0) y usar el despachador devuelto para llamar a sus funciones asincrónicas. La función devuelve una vista sobre el despachador original que limita el paralelismo a un límite que usted proporcione. Puedes usarlo así:
val limit = 2 // Or some other number val dispatcher = Dispatchers.Default val limitedDispatcher = dispatcher.limitedParallelism(limit) for (n in 0..100) { scope.launch(limitedDispatcher) { executeTask(n) } }Su pregunta, tal como se hizo, requiere la respuesta de @marstran. Si lo que desea es que no se ejecuten activamente más de N corrutinas en un momento dado (en paralelo), entonces el limitedParallelism es el camino a seguir:
val maxThreads: Int = TODO("some max number of threads") val limitedDispatcher = Dispatchers.Default.limitedParallelism(maxThreads) elements.forEach { elt -> scope.launch(limitedDispatcher) { doSomething(elt) } }Ahora, si lo que desea es incluso limitar la concurrencia , de modo que, como máximo, se ejecuten N corrutinas al mismo tiempo (potencialmente entrelazadas), independientemente de los subprocesos, podría usar un semáforo en su lugar:
val maxConcurrency: Int = TODO("some max number of concurrency coroutines") val semaphore = Semaphore(maxConcurrency) elements.forEach { elt -> scope.async { semaphore.withPermit { doSomething(elt) } } }También puede combinar ambos enfoques.
Otras respuestas ya explicaron que depende de si necesita limitar el paralelismo o la concurrencia. Si necesita limitar la concurrencia, puede hacerlo de manera similar a su solución original, pero con un solo canal:
val channel = Channel<T>() repeat(n) { launch { for(t in channel){ fn(t) } } } También tenga en cuenta que la offer() en su ejemplo no garantiza que la tarea se ejecute alguna vez. Si el siguiente consumidor en el round robin todavía está ocupado con la tarea anterior, simplemente se ignora la nueva tarea.