NodeJS实现AWS SQS订阅器:处理任务时跳过消息拉取
解决NodeJS中SQS轮询与异步任务冲突的问题
嘿,这个场景我太熟了!你想要的就是在视频转码这个耗时异步任务没做完的时候,暂停SQS的消息拉取,对吧?先帮你看看你测试代码里的问题,再给你两个靠谱的实现方案。
你的测试代码问题出在哪?
你尝试用asyncthingy.promise来判断异步任务是否完成,但这里有两个关键错误:
- 你的
asyncCall()函数虽然是async函数,但判断逻辑用错了——Promise对象本身并没有promise这个属性,所以!asyncthingy.promise这个判断永远为true,完全起不到“跳过”的作用; - 没有用明确的状态变量跟踪任务是否在运行,导致每次定时器触发都会重新调用异步任务。
修复后的测试代码可以这样写,用一个状态变量来控制:
let isProcessing = false; function resolveAfter4Seconds() { return new Promise(resolve => { setTimeout(() => { resolve('resolved'); }, 4000); }); } async function asyncCall() { console.log("I'm the async call"); const result = await resolveAfter4Seconds(); console.log(result); } const testTimer = async () => { if (isProcessing) { console.log("当前有异步任务在处理,跳过本次执行"); return; } isProcessing = true; await asyncCall(); isProcessing = false; } setInterval(testTimer, 1500);
针对SQS+视频转码的可行方案
方案一:状态变量 + setInterval
用布尔变量跟踪转码任务的运行状态,每次定时器触发时先检查状态,再决定是否拉取SQS消息。这种方式适合需要严格按照固定间隔(比如每20秒)尝试拉取的场景。
const AWS = require('aws-sdk'); const sqs = new AWS.SQS({ region: '你的区域' }); let isProcessing = false; const POLL_INTERVAL = 20000; // 20秒 const QUEUE_URL = '你的SQS队列URL'; // 处理SQS消息的核心函数 async function processMessage(message) { try { console.log(`开始处理消息: ${message.MessageId}`); // 执行视频转码操作(替换成你的实际转码逻辑) await transcodeVideo(message.Body); // 转码完成后删除消息(根据业务需求可选) await sqs.deleteMessage({ QueueUrl: QUEUE_URL, ReceiptHandle: message.ReceiptHandle }).promise(); console.log(`消息处理完成: ${message.MessageId}`); } catch (err) { console.error(`处理消息失败: ${err.message}`); // 这里可以根据需求添加重试逻辑或死信队列处理 } } // 模拟视频转码函数(最长5分钟) async function transcodeVideo(videoInfo) { return new Promise(resolve => { setTimeout(() => { console.log('视频转码完成'); resolve(); }, 300000); // 5分钟的模拟时长 }); } // 启动轮询 function startPolling() { setInterval(async () => { if (isProcessing) { console.log('当前有转码任务在运行,跳过本次SQS拉取'); return; } isProcessing = true; try { // 长轮询拉取SQS消息 const response = await sqs.receiveMessage({ QueueUrl: QUEUE_URL, MaxNumberOfMessages: 1, WaitTimeSeconds: 20 // 长轮询的等待时间 }).promise(); if (response.Messages && response.Messages.length > 0) { await processMessage(response.Messages[0]); } else { console.log('没有新消息,等待下一次轮询'); } } catch (err) { console.error('拉取SQS消息失败:', err.message); } finally { // 无论成功失败,都重置处理状态 isProcessing = false; } }, POLL_INTERVAL); } startPolling();
方案二:递归调度(更可靠的方式)
放弃setInterval,改用“任务完成后再调度下一次轮询”的递归方式。这种方式不会出现“任务未完成但定时器已触发”的情况,间隔是任务完成后等待20秒再拉取,更适合转码时长不稳定的场景。
const AWS = require('aws-sdk'); const sqs = new AWS.SQS({ region: '你的区域' }); const POLL_INTERVAL = 20000; // 20秒 const QUEUE_URL = '你的SQS队列URL'; async function pollAndProcess() { try { // 长轮询拉取消息 const response = await sqs.receiveMessage({ QueueUrl: QUEUE_URL, MaxNumberOfMessages: 1, WaitTimeSeconds: 20 }).promise(); if (response.Messages && response.Messages.length > 0) { const message = response.Messages[0]; console.log(`开始处理消息: ${message.MessageId}`); await transcodeVideo(message.Body); await sqs.deleteMessage({ QueueUrl: QUEUE_URL, ReceiptHandle: message.ReceiptHandle }).promise(); console.log(`消息处理完成: ${message.MessageId}`); } else { console.log('没有新消息'); } } catch (err) { console.error('轮询/处理消息出错:', err.message); } finally { // 无论结果如何,等待指定时间后再进行下一次轮询 setTimeout(pollAndProcess, POLL_INTERVAL); } } // 模拟转码函数 async function transcodeVideo(videoInfo) { return new Promise(resolve => { setTimeout(resolve, 300000); }); } // 启动轮询 pollAndProcess();
方案对比
- 方案一:严格按照固定间隔尝试拉取,适合对轮询时间精度要求高的场景,但如果转码时间超过20秒,会自动跳过后续的拉取直到任务完成。
- 方案二:任务完成后才开始计时下一次轮询,完全避免了任务重叠的可能,逻辑更简洁可靠,是我更推荐的方式。
内容的提问来源于stack exchange,提问作者Shaidar
相关产品推荐
相关产品推荐

