【问题标题】:Why does Promise.race not resolve in kafkajs eachMessage callback为什么 Promise.race 不能在 kafkajs eachMessage 回调中解析
【发布时间】:2021-05-15 09:15:29
【问题描述】:

我已经定义了一个这样的承诺......

    const result = await Promise.race([
      new Promise(resolve => {
        consumer.run({
          eachMessage: ({ message }) => {
            const data = JSON.parse(message.value.toString());
            if (data.payload.template
              && data.payload.template.id === '...'
              && data.payload.to[0].email === email) {
              console.log('Should resolve!')
              resolve(data.payload.template.variables.link);
              console.log('resolved');

              consumer.pause();
              consumer.disconnect();
            }
          },
        });
      }),
      new Promise((_, reject) => setTimeout(reject, 3000))
    ]);
    console.log('result is ', result);
    return result;

我可以解决,但最后没有打印结果,似乎超时和实际承诺都没有按预期工作?这是为什么?我怀疑这与在 kafka js 回调中使用 resolve 有关?


更新:似乎它的Promise.race() 没有解决,但为什么呢?

【问题讨论】:

  • " 似乎它的 Promise.race() 没有解决" - 它是否在拒绝并且你正在以某种方式默默地吞下错误?
  • 不相关的观察:您可能应该将consumer.pause()consumer.disconnect() 移动到超时承诺处理程序中。这样,无论如何,消费者最终都会被暂停和断开连接。在您当前的实现中,它只会在成功的情况下暂停和断开连接。 (这可能是故意的,也可能不是故意的,我只是注意到了。)
  • 我的错...@Tomalak 你是对的...
  • 我也会按照建议删除暂停和断开连接

标签: node.js promise kafkajs


【解决方案1】:

我怀疑您的“成功方面”承诺无意中抛出,而您正在默默地吞下错误。

使用 consumer 的模型最小实现(成功或失败 50/50),以下代码有效。

运行代码示例几次以查看这两种情况。

var consumer = {
  interval: null,
  counter: 0,
  run: function (config) {
    this.interval = setInterval(() => {
      this.counter++;
      console.log(`Consumer: message #${this.counter}`);
      config.eachMessage({message: this.counter});
    }, 250);
  },
  pause: function () {
    console.log('Consumer: paused');
    clearInterval(this.interval);
  },
  disconnect: function () {
    console.log('Consumer: disconnected');    
  }
};

Promise.race([
  new Promise(resolve => {
    const expectedMsg = Math.random() < 0.5 ? 3 : 4;
    consumer.run({
      eachMessage: ({ message }) => {
        if (message === expectedMsg) resolve("success");
      }
    });
  }),
  new Promise((_, reject) => setTimeout(() => {
    reject('timeout');
    consumer.pause();
    consumer.disconnect();
  }, 1000))
]).then((result) => {
  console.log(`Result: ${result}`);
}).catch((err) => {
  console.log(`ERROR: ${err}`);
});

我还将consumer.pause()consumer.disconnect() 移至“超时端”承诺,这样消费者就可以保证断开连接,尽管在成功的情况下它可能运行的时间可能比必要的时间长一点。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-04-27
    • 2011-06-11
    • 2020-11-16
    • 2012-06-18
    • 1970-01-01
    • 2020-08-12
    • 1970-01-01
    • 2010-09-19
    相关资源
    最近更新 更多