【问题标题】:Node.js EventEmitter events not sharing event loopNode.js EventEmitter 事件不共享事件循环
【发布时间】:2015-07-14 03:55:31
【问题描述】:

也许根本问题是我正在使用的node-kafka 模块是如何实现的,但也许不是,所以我们开始...

使用 node-kafa 库时,我遇到了订阅 consumer.on('message') 事件的问题。该库使用标准的events 模块,所以我认为这个问题可能足够通用。

我的实际代码结构又大又复杂,所以这里是一个基本布局的伪示例来突出我的问题。 (注意:这段代码 sn-p 未经测试,所以我可能在这里有错误,但无论如何这里的语法没有问题)

var messageCount = 0;
var queryCount = 0;

// Getting messages via some event Emitter
consumer.on('message', function(message) {
    message++;
    console.log('Message #' + message);

    // Making a database call for each message
    mysql.query('SELECT "test" AS testQuery', function(err, rows, fields) {
        queryCount++;
        console.log('Query   #' + queryCount);
    });
})

我在这里看到的是,当我启动服务器时,有 100,000 条左右的积压消息是 kafka 想要给我的,它通过事件发射器这样做。所以我开始收到消息。获取和记录所有消息大约需要 15 秒。

假设 mysql 查询相当快,这是我期望看到的输出:

Message #1
Message #2
Message #3
...
Message #500
Query   #1
Message #501
Message #502
Query   #2
... and so on in some intermingled fashion

我希望这是因为我的第一个 mysql 结果应该很快准备好,并且我希望结果在事件循环中轮流处理响应。我实际得到的是:

Message #1
Message #2
...
Message #100000
Query   #1
Query   #2
...
Query   #100000

在能够处理 mysql 响应之前,我会收到每条消息。所以我的问题是,为什么?为什么在所有消息事件都完成之前,我无法获得单个数据库结果?

另一个注意事项:我在 node-kafka 中的 .emit('message') 和我的代码中的 mysql.query() 设置了一个断点,并且我正在轮流击中它们。因此,在进入我的事件订阅者之前,似乎所有 100,000 个发射都没有预先堆积。所以我就这个问题提出了第一个假设。

非常感谢您的想法和知识:)

【问题讨论】:

  • 如果将存储的消息数量增加到更大的数量会怎样?是不是你的 mysql 就是这么慢?
  • @Avery 我很想知道,但是当我只用一条要处理的消息复制它时,我什至无法感知 mysql 响应的延迟。这也都在本地运行。而且实际的 mysql 查询非常简单(只是从单个表行中选择约 8 个字段,而该表现在只有大约 60 行)
  • 如果这个例子实际上代表了你的代码,那么我也迷路了。你真的能用这个例子产生这个结果吗?我没有可用于测试的 MySQL 实例。
  • 您是否为node-kafka 配置了足够大的fetchMaxBytes 值,以便在一个请求中传输所有这100K 条消息? EventEmitter 是同步的,它不使用 Node 事件循环,因此如果一次有 100K 消息进入,它们可能会在你的异步代码有机会运行之前全部发出。
  • @robertklep 谢谢!因此,在 kafka-node 的示例中,他们使用fetchMaxBytes: 1024*10 显示了默认覆盖示例。在其他默认覆盖中,它们的值等于默认值,他们甚至注意到了这一点,所以我认为这个属性就是这种情况。您的问题启发了我查看他们的代码并查看其默认值实际上是fetchMaxBytes: 1024*1024。所以是的,我实际上是在一个请求中接收所有消息。而且我不知道 EventEmitter 是同步的 :)

标签: javascript node.js eventemitter


【解决方案1】:

node-kafka 驱动程序使用了相当宽泛的缓冲区大小 (1M),这意味着它将从 Kafka 获取尽可能多的消息,这些消息将适合缓冲区。如果服务器积压,并且根据消息大小,这可能意味着(数万)条消息随一个请求进入。

因为 EventEmitter 是同步的(它不使用 Node 事件循环),这意味着驱动程序将向其侦听器发出(成千上万)个事件,并且由于它是同步的,它不会屈服于 Node事件循环,直到所有消息都已传递。

我认为您无法解决事件交付的泛滥问题,但我不认为具体的事件交付存在问题。更可能的问题是为每个事件启动异步操作(在本例中为 MySQL 查询),这可能会使数据库充满查询。

一种可能的解决方法是使用队列而不是直接从事件处理程序执行查询。例如,使用async.queue,您可以限制并发(异步)任务的数量。队列的“worker”部分将执行 MySQL 查询,而在事件处理程序中,您只需将消息推送到队列中。

【讨论】:

  • 谢谢@robertklep。我会试试 async.queue。我正在处理自己的队列,因此只有一个 mysql 查询和内存缓存结果以供等待请求使用,但我怀疑维护/测试良好的模块会更好:)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-05-23
  • 1970-01-01
  • 2015-10-13
  • 2018-01-15
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多