如何在嵌套async/await执行完成后执行指定函数(文件文本提取与AMQP消息队列场景)
解决文本提取与消息发送完成后关闭AMQP连接的问题
我明白你的核心痛点:当前代码里textract.fromFileWithPath的回调式异步逻辑和外层的async/await没有正确协同,导致startExtraction会在所有文件的文本提取、数据库存储和消息队列发布完成前就执行完毕,没法准确触发AMQP连接的关闭操作。
问题根源
textract.fromFileWithPath是基于回调的异步API,你在for循环中调用它时,并不会等待回调内的逻辑(文本提取、数据库保存、消息发送)完成就会进入下一次循环。这导致startExtraction里的console.log("Done")会在所有文件处理完成前就打印,后续的连接关闭操作自然也会提前执行,而此时还有大量异步任务在后台运行。
解决方案步骤
要解决这个问题,我们需要把回调式API转换成Promise,再用Promise.all等待所有文件的处理任务全部完成:
1. 包装textract API为Promise
首先把textract.fromFileWithPath封装成返回Promise的函数,让它能和async/await兼容:
// 把回调式的textract API包装成Promise const textractFromFile = (filePath) => { return new Promise((resolve, reject) => { textract.fromFileWithPath(filePath, (err, text) => { if (err) reject(err); else resolve(text); }); }); };
2. 修改startExtraction函数
重构函数逻辑,用Promise.all等待所有文件的处理任务完成:
const startExtraction = async (dir, channel) => { console.log("Started"); const files = fs.readdirSync(dir); // 把每个文件的处理流程转换成Promise const processingTasks = files.map(async (file) => { const native = `${root}\\${file}`; try { const text = await textractFromFile(native); // 注意:这里的payload需要包含text等数据库需要的字段,比如 payload = { text, ... } const doc = await File.create(payload); channel.sendToQueue(queue.MLIFY, Buffer.from(JSON.stringify(doc))); console.log("Sent "+ doc._id); } catch(err) { console.log(`处理文件${file}出错:`, err); } }); // 等待所有文件的处理任务全部完成 await Promise.all(processingTasks); console.log("所有文件处理完成"); };
3. 修改initiateExtraction函数
现在startExtraction会返回一个Promise,只有当所有任务完成后才会resolve,你可以直接在它之后安全关闭连接:
const initiateExtraction = async (job) => { let conn; try { conn = await amqp.connect('amqp://localhost'); const channel = await conn.createChannel(); await channel.assertQueue(queue.MLIFY, { durable: true }); // 等待所有文本提取、数据库存储、消息发送任务完成 await startExtraction(job, channel); console.log("所有任务已完成,准备关闭AMQP连接"); // 先关闭channel再关闭连接 await channel.close(); await conn.close(); } catch (error) { console.error("初始化提取过程出错:", error); // 即使出错也要尝试关闭连接,避免资源泄漏 if (conn) { try { await conn.close(); } catch (closeErr) { console.error("关闭连接时出错:", closeErr); } } } };
关键说明
Promise.all是核心:它会等待所有传入的Promise任务都完成后才会继续执行后续代码,确保我们不会提前关闭连接。- 回调转Promise:这是让
async/await能正确管理异步流程的前提,把分散的回调逻辑统一纳入到Promise链中。 - 错误处理优化:在
initiateExtraction中增加了连接的兜底关闭逻辑,避免因过程出错导致连接资源泄漏。
内容的提问来源于stack exchange,提问作者Venkatesh Dharavath
相关产品推荐
相关产品推荐

