我可以限制kafka节点消费者的消费吗?

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-node,它似乎与Java版本相比具有有限的api

Mat*_*Sax 5

在 Kafka 中,轮询和处理应该以协调/同步的方式发生。即,在每次轮询之后,您应该先处理所有接收到的数据,然后再进行下一次轮询。此模式将自动将消息数量限制为您的客户端可以处理的最大吞吐量。

像这样(伪代码):

while(isRunning) {
  messages = poll(...)
  for(m : messages) {
    process(m);
  }
}
Run Code Online (Sandbox Code Playgroud)

(这就是为什么没有参数“fetch.max.messages”的原因——你只是不需要它。)


Nik*_*ain 5

我有类似的情况,我正在消费来自Kafka的消息,不得不限制消费,因为我的消费者服务依赖于有自己约束的第三方API.

我使用async/queue了一个async/cargo叫做asyncTimedCargo批处理的包装器.货物从kafka-consumer获取所有消息,并在达到大小限制batch_config.batch_size或超时时将其发送到队列batch_config.batch_timeout. async/queue提供saturatedunsaturated回调,如果队列任务工作人员忙,您可以使用它们来停止消耗.这将阻止货物填满,您的应用程序不会耗尽内存.消费将在不满足时恢复.

//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)


Tob*_*iSH 0

据我所知,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)