Mat*_*ner 5 node.js apache-kafka
看起来像我的kafka节点消费者:
var kafka = require('kafka-node');
var consumer = new Consumer(client, [], {
...
});
Run Code Online (Sandbox Code Playgroud)
在某些情况下,获取的消息太多了.有没有办法限制它(例如,每秒接受不超过1000条消息,可能使用暂停api?)
在 Kafka 中,轮询和处理应该以协调/同步的方式发生。即,在每次轮询之后,您应该先处理所有接收到的数据,然后再进行下一次轮询。此模式将自动将消息数量限制为您的客户端可以处理的最大吞吐量。
像这样(伪代码):
while(isRunning) {
messages = poll(...)
for(m : messages) {
process(m);
}
}
Run Code Online (Sandbox Code Playgroud)
(这就是为什么没有参数“fetch.max.messages”的原因——你只是不需要它。)
我有类似的情况,我正在消费来自Kafka的消息,不得不限制消费,因为我的消费者服务依赖于有自己约束的第三方API.
我使用async/queue了一个async/cargo叫做asyncTimedCargo批处理的包装器.货物从kafka-consumer获取所有消息,并在达到大小限制batch_config.batch_size或超时时将其发送到队列batch_config.batch_timeout.
async/queue提供saturated和unsaturated回调,如果队列任务工作人员忙,您可以使用它们来停止消耗.这将阻止货物填满,您的应用程序不会耗尽内存.消费将在不满足时恢复.
//cargo-service.js
module.exports = function(key){
return new asyncTimedCargo(function(tasks, callback) {
var length = tasks.length;
var postBody = [];
for(var i=0;i<length;i++){
var message ={};
var task = JSON.parse(tasks[i].value);
message = task;
postBody.push(message);
}
var postJson = {
"json": {"request":postBody}
};
sms_queue.push(postJson);
callback();
}, batch_config.batch_size, batch_config.batch_timeout)
};
//kafka-consumer.js
cargo = cargo-service()
consumer.on('message', function (message) {
if(message && message.value && utils.isValidJsonString(message.value)) {
var msgObject = JSON.parse(message.value);
cargo.push(message);
}
else {
logger.error('Invalid JSON Message');
}
});
// sms-queue.js
var sms_queue = queue(
retryable({
times: queue_config.num_retries,
errorFilter: function (err) {
logger.info("inside retry");
console.log(err);
if (err) {
return true;
}
else {
return false;
}
}
}, function (task, callback) {
// your worker task for queue
callback()
}), queue_config.queue_worker_threads);
sms_queue.saturated = function() {
consumer.pause();
logger.warn('Queue saturated Consumption paused: ' + sms_queue.running());
};
sms_queue.unsaturated = function() {
consumer.resume();
logger.info('Queue unsaturated Consumption resumed: ' + sms_queue.running());
};
Run Code Online (Sandbox Code Playgroud)
据我所知,API 没有任何类型的限制。但是两个消费者(Consumer 和 HighLevelConsumer)都有一个“pause()”函数。因此,如果您收到太多消息,您可以停止消费。也许这已经提供了您所需要的。
请记住正在发生的事情。您向代理发送获取请求并获取一批消息。您可以配置要获取的消息的最小和最大大小(根据文档而不是消息数):
{
....
// This is the minimum number of bytes of messages that must be available to give a response, default 1 byte
fetchMinBytes: 1,
// The maximum bytes to include in the message set for this partition. This helps bound the size of the response.
fetchMaxBytes: 1024 * 1024,
}
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
4126 次 |
| 最近记录: |