Estoy tratando de ejecutar una serie de promesas en paralelo en un lote, pero comienza el siguiente lote antes de que haya terminado de procesar el primer lote. El código funciona cuando uso promesas de tiempo de espera simples. ¿Puede alguien ayudarme a darme un poco de ayuda?
await asyncFunctionsInBatches.reduce(async (previousBatch, currentBatch, index) => { await previousBatch; console.time(`batch ${index}`); console.log(`Processing batch ${index}...`); const currentBatchPromises = currentBatch.map(async (row) => await sqlFetch(queryDate, queryDate1, username, password, row) ); await Promise.all(currentBatchPromises); console.timeEnd(`batch ${index}`); }, Promise.resolve()); }Esta es la promesa que estoy usando
const sqlFetch = async ( queryDate, queryDate1, username, password, rowNumber ) => { ( await createUnixSocketPool(username, password)).query( fetchQuery.sql, [queryDate, queryDate1, rowNumber], async function (err, result) { if (err) { console.log(err); } console.log(result) return publishMessage(result).catch(console.error); } ); };En sqlFetch , está esperando el resultado de createUnixSocketPool pero no el resultado de la query :
const sqlFetch = async ( queryDate, queryDate1, username, password, rowNumber ) => { ( await createUnixSocketPool(username, password)).query( // −−−−− −−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−−− −−−−−− // ↑ ↑ ↑ // \−−− only applies to / not to −−−−−−−−−−−−−−−/ fetchQuery.sql, [queryDate, queryDate1, rowNumber], async function (err, result) { if (err) { console.log(err); } console.log(result) return publishMessage(result).catch(console.error); } ); }; Como resultado, la promesa de sqlFetch se cumple de inmediato con undefined , mientras la consulta aún se está ejecutando.
Sin saber qué espera la query de su devolución de llamada (le está dando una promesa; ¿realmente espera una?) o si devuelve una promesa, es difícil ayudar, aparte de decir que si devuelve una promesa, await eso, y si no es así, maneje eso (por ejemplo, como se describe en las respuestas a esta pregunta ).
Aquí hay una actualización basada en la información obtenida de la respuesta que publicaste:
const sqlFetch = ( queryDate, queryDate1, username, password, rowNumber ) => new Promise((resolve, reject) => { createUnixSocketPool(username, password) .then(pool => { pool.query( fetchQuery.sql, [queryDate, queryDate1, rowNumber], (err, result) => { if (err) { reject(err); } else { resolve(publishMessage(result)); } } ); }) .catch(reject); }); O mejor aún, escriba un contenedor habilitado para promesas para la query como se describe en el enlace anterior :
function asyncQuery(pool, ...args) { return new Promise((resolve, reject) => { pool.query(...args, (err, result) => { if (err) { reject(err); } else { resolve(result); } }); }); }y luego es realmente agradable y simple:
const sqlFetch = async ( queryDate, queryDate1, username, password, rowNumber ) => { const pool = await createUnixSocketPool(username, password) const result = await asyncQuery(pool, fetchQuery.sql, [queryDate, queryDate1, rowNumber]); return publishMessage(result) });Agregué una promesa
const sqlFetch = ( queryDate, queryDate1, username, password, rowNumber ) => new Promise(async (resolve, reject) => { (await createUnixSocketPool(username, password) .catch(reject)).query( fetchQuery.sql, [queryDate, queryDate1, rowNumber], function (err, result) { if (err) reject(err); else publishMessage(result).catch(reject).then(resolve); } ); });