如何监听Agenda作业完成以获取加密货币提币TXID?
监听Agenda单个作业完成事件获取TXID的解决方案
我来帮你搞定这个问题!Agenda本身确实没有直接提供单个作业的完成监听API,但我们有几个优雅的实现方式,完全不用setInterval或者检查transactionStepCompleted:
方法1:全局事件+作业ID过滤
Agenda有全局的success和fail事件,我们可以在创建作业时记住它的ID,然后在全局事件里过滤出目标作业,拿到更新后的TXID。
修改你的作业创建代码:
const job = global.agenda.create('withdraw_order', { userId: user._id, recipientAddress: addr, amount, txid: null }); job.save(async (err) => { if (err) return false; // 只处理当前创建的这个作业 const handleJobCompletion = async (completedJob) => { if (completedJob.attrs._id.toString() === job.attrs._id.toString()) { // 这里就能拿到作业完成后更新的txid了 msg.reply(`Successfully withdrawn! TXID: ${completedJob.attrs.data.txid}`); // 移除监听,防止内存泄漏 agenda.removeListener('success', handleJobCompletion); agenda.removeListener('fail', handleJobCompletion); } }; // 监听作业成功和失败事件 agenda.on('success', handleJobCompletion); agenda.on('fail', handleJobCompletion); });
这个思路很简单:利用Agenda自带的全局事件,通过作业ID精准匹配我们创建的那个作业,事件触发时就能拿到完成后的完整作业数据,包括更新后的TXID。
方法2:内存事件发射器(适合单实例应用)
如果你的应用是单实例部署,用内存EventEmitter来传递作业完成通知会更直接:
首先,在应用初始化时创建一个全局的EventEmitter:
const EventEmitter = require('events'); global.jobNotificationEmitter = new EventEmitter();
然后修改作业创建逻辑:
const job = global.agenda.create('withdraw_order', { userId: user._id, recipientAddress: addr, amount, txid: null }); // 用作业ID作为事件名,确保唯一 const jobCompletedEvent = job.attrs._id.toString(); // 只监听一次这个作业的完成事件 jobNotificationEmitter.once(jobCompletedEvent, (txid, error) => { if (error) { msg.reply(`Withdrawal failed: ${error.message}`); } else { msg.reply(`Successfully withdrawn! TXID: ${txid}`); } }); job.save((err) => { if (err) { // 保存失败直接触发通知 jobNotificationEmitter.emit(jobCompletedEvent, null, err); return false; } });
接下来修改你的作业定义,在作业完成时触发事件:
agenda.define('withdraw_order', async function(job, done) { try { // 注意这里要让performWithdraw返回最终的txid const txid = await paymentProcessor.performWithdraw(job); // 触发通知,传递TXID global.jobNotificationEmitter.emit(job.attrs._id.toString(), txid, null); done(); } catch(e) { // 失败时传递错误信息 global.jobNotificationEmitter.emit(job.attrs._id.toString(), null, e); done(e); } });
还要调整performWithdraw函数,让它返回TXID而不是包装对象:
async performWithdraw(options) { try { // withdraw函数已经会返回sendID(也就是TXID) const txid = await this.withdraw(options); return txid; } catch(e) { this.reportException(e); throw e; // 抛出错误,让上层作业定义处理 } }
这样一来,作业完成后会立刻通过EventEmitter通知到创建作业的地方,完美拿到TXID。
方法3:MongoDB Change Stream(适合分布式场景)
因为Agenda底层用MongoDB存储作业,我们可以直接利用MongoDB的Change Stream来监听单个作业文档的更新:
修改作业创建代码:
const job = global.agenda.create('withdraw_order', { userId: user._id, recipientAddress: addr, amount, txid: null }); job.save(async (err) => { if (err) return false; // 获取Agenda的MongoDB集合 const jobCollection = agenda._collection; // 只监听当前作业的completed字段变为true的更新 const changeStream = jobCollection.watch([ { $match: { 'documentKey._id': job.attrs._id, 'updateDescription.updatedFields.completed': true }} ]); // 监听变更事件 changeStream.on('change', async () => { // 查询最新的作业数据 const completedJob = await agenda.jobs({ _id: job.attrs._id }); msg.reply(`Successfully withdrawn! TXID: ${completedJob[0].attrs.data.txid}`); // 关闭流,防止内存泄漏 changeStream.close(); }); });
这个方法的好处是支持分布式部署,只要所有实例连接同一个MongoDB,就能正确收到作业完成的通知。
关键注意点
- 不管用哪种方法,都一定要记得移除监听或者关闭流,避免内存泄漏。
- 如果是分布式应用,内存EventEmitter就不适用了,这时候可以用Redis Pub/Sub或者类似的分布式消息队列替代。
内容的提问来源于stack exchange,提问作者user9761835
相关产品推荐
相关产品推荐

