En el código fuente de mi programa tengo la siguiente función (función de limitación de concurrencia Promise, similar a pLimit):
async function promiseMapLimit( array, poolLimit, iteratorFn, ) { const ret = []; const executing = []; for (const item of array) { const p = Promise.resolve().then(() => iteratorFn(item, array)); ret.push(p); if (poolLimit <= array.length) { const e = p.then(() => executing.splice(executing.indexOf(e), 1)); executing.push(e); if (executing.length >= poolLimit) { await Promise.race(executing); } } } return Promise.all(ret); }Funciona correctamente, por lo que si le paso una matriz de números [1..99] y trato de multiplicarlo por 2, dará el resultado correcto [0..198].
const testArray = Array.from(Array(100).keys()); promiseMapLimit(testArray, 20, async (value) => value * 2).then((result) => console.log(result) );Ejemplo de código: js playground .
Pero no puedo entender su lógica, durante la depuración noté que agrega promesas en trozos de 20 y solo después de eso va más allá: 
Por ejemplo, este bloque de código:
for (const item of array) { const p = Promise.resolve().then(() => iteratorFn(item, array)); ret.push(p);iterará sobre 20 elementos de una matriz ( ¿por qué no los 100? )
aquí igual:
if (poolLimit <= array.length) { const e = p.then(() => executing.splice(executing.indexOf(e), 1)); executing.push(e); agregará solo 20 elementos a la matriz en executing y solo después de ese paso dentro del bloque de código if (executing.length >= poolLimit) .
Estaría muy agradecido por la explicación de cómo funciona esta función.
Pregunta muy interesante! Creo que la parte importante del código aquí es Promise.race(...) que se resuelve tan pronto como se resuelve una de las promesas.
He agregado una función de sleep con un factor aleatorio (hasta 6 segundos) para visualizar mejor cómo funciona esto.
La funcionalidad esperada es la siguiente: siempre queremos que se ejecuten 20 promesas en paralelo, y una vez que una termine, se ejecutará la siguiente en la cola.
De forma visual, se vería así, para un límite de 3 y 10 promesas: en el siguiente ejemplo, puede notar que en cada momento hay 3 promesas activas (excepto cuando terminan):
PromiseID | Start End | 0 [====] 1 [==] 2 [======] 3 [==========] 4 [====] 5 [================] 6 [==] 7 [====] 8 [======] 9 [========]El código para crear el retraso aleatorio se encuentra a continuación:
// Create the utility sleep function const sleep = x => new Promise(res => setTimeout(res, x)) async function promiseMapLimit(array, poolLimit, iteratorFn) { const ret = []; const executing = []; for (const item of array) { const p = Promise.resolve().then(() => iteratorFn(item, array)); ret.push(p); console.log(ret.length) if (poolLimit <= array.length) { const e = p.then(() => executing.splice(executing.indexOf(e), 1)); executing.push(e); if (executing.length >= poolLimit) { console.log(`Running batch of ${executing.length} promises.`); await Promise.race(executing); // As ssoon one of the promise finishes, we continue the loop. console.log("Resolved one promise.") } } } return Promise.all(ret); } const testArray = Array.from(Array(100).keys()); promiseMapLimit(testArray, 20, async (value) => { // Log console.log(`Computing iterator fn for ${value}`) await sleep(3000 + Math.random() * 3000); return value * 2 }).then((result) => console.log(result) );iterará sobre 20 elementos de una matriz (¿por qué no los 100?)
Al inicio, como en el gráfico, no iterará los 100 elementos, sino los primeros 20 elementos y luego el ciclo se pausará con await Promise.race(...) (porque executing.length >= poolLimit será verdadero después de iterar 20 artículos).
Una vez que se completa una promesa, se eliminará de la matriz de executing executing.splice(executing.indexOf(e), 1) .
Creo que las cosas se vuelven más claras cuando hay un retraso ( await sleep(...) ) para que podamos simular una operación asíncrona real (como una solicitud de base de datos, etc.).
Por favor, avíseme si hay algo más que no esté claro.
Tienes await dentro de la función asíncrona. Esto funciona más o menos de la siguiente manera:
await palabra claveawait En su caso, itera 20 veces, luego detiene todo una vez que alcanza un límite. Luego, una vez que se resuelve al menos una promesa dentro de ret , continúa.
Lo siguiente que sucede es que una vez que se resuelve alguna de las promesas, se elimina de la matriz. Pero dado que casi todo sucede instantáneamente, verá que resuelve las 20 promesas, se llena con otras 20. Si hace que su iteratorFn sea más lento con retrasos aleatorios, verá que ese grupo se llena constantemente hasta 20 y luego reemplaza casi de inmediato espacio liberado en la piscina con nueva promesa, mientras que quedan al menos algunos elementos.
Reemplacemos su iteratorFn con esto y llámelo:
let iter = async (value) => { // randomly delay each calculation to 1, 2 or 3 seconds return new Promise(resolve => setTimeout(resolve, [1000, 2000, 3000][Math.floor(Math.random() * 3)], value * 2)) } promiseMapLimit(testArray, 20, iter).then((result) => console.log(result) ); Y registremos la cantidad de elementos dentro de la executing una vez que se resuelva una promesa:
if (poolLimit <= array.length) { const e = p.then(() => { executing.splice(executing.indexOf(e), 1); // logging what is left console.log({l: executing.length}) }); executing.push(e); if (executing.length >= poolLimit) { await Promise.race(executing); } } De esta manera, en la consola, verá que el registro comienza con {l: 19} , ya que el grupo se llena y luego se resuelve una promesa. Y continuará, hasta el final, donde el registro pasará de 19 a 0.
Será divertido escribir esta respuesta.
async function promiseMapLimit( array, poolLimit, iteratorFn, ) { const ret = []; const executing = []; for (const item of array) { const p = Promise.resolve().then(() => iteratorFn(item, array)); ret.push(p); if (poolLimit <= array.length) { const e = p.then(() => executing.splice(executing.indexOf(e), 1)); executing.push(e); if (executing.length >= poolLimit) { await Promise.race(executing); } } } return Promise.all(ret); }Así que hay tres cosas sucediendo en este fragmento de código.
Promise.raceIgnore 2 por algún tiempo y analicemos 1 y 3.
Veamos un escenario simple.
promiseMapLimit([1, 2, 3], 1, async function (i) => { return 2 * i }); Esta línea llama a la función anterior. ahora debido a esto como const p = Promise.resolve().then(() => iteratorFn(item, array)); ret ahora tiene [Promesa(2), Promesa(4), Promesa(6)]
Ahora, debido a que poolLimit es menor que array.length, ingresamos al bloque if.
const e = p.then(() => executing.splice(executing.indexOf(e), 1)); esta línea agrega un bloque then que se ejecutará después de que se resuelvan los elementos de la promesa de la matriz ret. En este bloque, cada promesa itera a través de la matriz de ejecución, se encuentra y se elimina.
es equivalente a:
const e = p.then(() => { 1. Keep e as closure variable for later, for when the promise does get resolved in Promise.race line 2. find that e in executing array 3. Remove it }); Ahora que comprende el 1 y el 3. Vayamos al 2. El bucle se ejecuta para array.length y cuando ejecuta.length >= poolLimit, espera la ejecución de las promesas en la cola de ejecución debido a que esta línea await Promise.race(executing);
Déjame reescribir la función para que sea más fácil de entender:
async function promiseMapLimit( array, poolLimit, iteratorFn, ) { const ret = []; const executing = []; for (const item of array) { const p = Promise.resolve().then(() => iteratorFn(item, array)); ret.push(p); } if (poolLimit > array.length) { // if poolLimit and greater that array just resolve all the promises and no need to worry of concurrency return Promise.all(ret); } for (const p of ret) { // Add self remove logic to each promise element const e = p.then(() => executing.splice(executing.indexOf(e), 1)); // Put the promise element in execution queue executing.push(e); // Whenever the execution ques is full wait for execution of the promises // Promise.race awaits for all the promises to finish execution // And because of the self removal logic execution queue also becomes empty if (executing.length >= poolLimit) { await Promise.race(executing); } } // Promise.all will take care of all the promises remaining in the last batch for which await Promise.race(executing); was not executed. return Promise.all(ret); }