Eri*_*son 6 javascript node.js eventemitter
也许潜在的问题是我使用的node-kafka模块是如何实现的,但也许不是,所以我们在这里......
使用node-kafa库,我遇到订阅consumer.on('message')事件的问题.该库正在使用标准events模块,所以我认为这个问题可能足够通用.
我的实际代码结构庞大而复杂,所以这里有一个基本布局的伪示例来突出我的问题.(注意:此代码段未经测试,因此我可能在此处有错误,但无论如何语法都没有问题)
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);
});
})
Run Code Online (Sandbox Code Playgroud)
我在这里看到的是,当我启动我的服务器时,卡夫卡希望给我的100,000条左右的积压消息,它通过事件发射器完成.所以我开始收到消息.获取并记录所有消息大约需要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
Run Code Online (Sandbox Code Playgroud)
我希望这是因为我的第一个mysql结果应该很快就准备好了,我希望结果可以在事件循环中轮流处理响应.我实际得到的是:
Message #1
Message #2
...
Message #100000
Query #1
Query #2
...
Query #100000
Run Code Online (Sandbox Code Playgroud)
在获得mysql响应之前,我收到了每条消息.所以我的问题是,为什么?为什么在所有消息事件完成之前我无法获得单个数据库结果?
另一个注意事项:我.emit('message')在node-kafka和mysql.query()我的代码中设置了一个断点,我将它们转为基于回合制.因此,在进入我的事件订阅者之前,似乎所有100,000个发射都没有预先堆叠.所以我对这个问题进行了第一次假设.
想法和知识将非常感激:)
该node-kafka驱动程序使用相当宽松的缓冲区大小(1M),这意味着它将从 Kafka 获取适合缓冲区的尽可能多的消息。如果服务器积压,并且根据消息大小,这可能意味着一个请求会传入(数万)万条消息。
因为 EventEmitter 是同步的(它不使用 Node 事件循环),这意味着驱动程序将向其侦听器发出(成千上万)个事件,并且由于它是同步的,因此它不会屈服于 Node 事件循环,直到所有消息均已送达。
我不认为你可以解决大量的事件传递,但我不认为具体的事件传递是有问题的。更可能的问题是为每个事件启动异步操作(在本例中为 MySQL 查询),这可能会导致数据库充满查询。
一种可能的解决方法是使用队列,而不是直接从事件处理程序执行查询。例如,async.queue您可以限制并发(异步)任务的数量。队列的“worker”部分将执行 MySQL 查询,而在事件处理程序中,您只需将消息推送到队列中。
| 归档时间: |
|
| 查看次数: |
431 次 |
| 最近记录: |