You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何监听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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 07:06:45