Tengo una secuencia grande que se genera con pereza y es demasiado grande para caber en la memoria.
Me gustaría procesar esta secuencia usando corrutinas para mejorar el rendimiento, en este ejemplo estoy usando 10 subprocesos paralelos para el procesamiento.
runBlocking(Dispatchers.IO.limitedParallelism(10)) { massiveLazySequenceOfThings.forEach { myThingToProcess -> print("I am processing $myThingToProcess") launch() { process(myThingToProcess) } } }El problema aquí es que la primera declaración de impresión se ejecutará para CADA elemento de la secuencia, por lo que para secuencias extremadamente grandes como la mía, esto causará OOM.
¿No hay forma de hacer que la iteración de mi secuencia sea "perezosa" en este ejemplo, de modo que solo se procese un número fijo en un momento dado?
¿Estoy obligado a usar canales aquí (¿posiblemente con un canal almacenado en búfer?) para forzar una llamada de bloqueo durante mi iteración de secuencia hasta que se procesen algunos elementos? O hay alguna otra solución más limpia que me falta.
En mi ejemplo actual, también estoy usando un supervisorScope para monitorear cada trabajo de procesamiento, por lo que, si es posible, también me gustaría conservarlo.
Puede limitar el número de operaciones paralelas con Semaphore .
runBlocking(Dispatchers.IO) { val semaphore = Semaphore(permits = 10) massiveLazySequenceOfThings.forEach { myThingToProcess -> semaphore.acquire() print("I am processing $myThingToProcess") launch { try { process(myThingToProcess) } finally { print("I am done processing $myThingToProcess") semaphore.release() } } } }El problema es probablemente que programa nuevas tareas más rápido de lo que es posible terminar las tareas en paralelo.
Puede superar este problema fragmentando sus datos en el tamaño que se puede procesar en paralelo (en su ejemplo, que es 10) y esperar a que termine cada fragmento antes de comenzar con el siguiente fragmento:
runBlocking(Dispatchers.IO) { massiveLazySequenceOfThings.chunked(10).forEach { chunk -> println("I am processing $chunk") chunk.map { async { process(it) } }.awaitAll() } }Ahora, cada uno de los 10 elementos de un fragmento se procesa en paralelo, pero cada fragmento en su conjunto se procesa secuencialmente, asegurándose de que no se quede sin memoria debido a demasiadas tareas programadas.
Una pequeña nota al margen: el resultado "I am done processing $myThingToProcess" en su ejemplo no es correcto. Es posible que el trabajo aún se esté ejecutando en su propio subproceso. Lo único que puede decir con certeza es que lo programó.