【发布时间】: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