2017-02-26 11 views
0

kafka-node ConsumerGroupを使用してトピックからメッセージを消費しています。 ConsumerGroupはメッセージを消費するときに外部APIを呼び出す必要があります。応答にはさらに時間がかかることがあります。 APIからの応答を得るまで、キューから次のメッセージの消費を制御して、メッセージが順次処理されるようにします。ConsumerGroupによってメッセージを処理する並行性を制御する方法

この動作をどのように制御する必要がありますか?

答えて

0

これは、我々は一度に1つのメッセージの処理を実装している方法です。

var async = require('async'); //npm install async 

//intialize a local worker queue with concurrency as 1 (only 1 event is processed at a time) 
var q = async.queue(function(message, cb) { 
      processMessage(message).then(function(ep) { 
      cb(); //this marks the completion of the processing by the worker 
     }); 
}, 1); 

// a callback function, invoked when queue is empty. 
q.drain = function() { 
    consumerGroup.resume(); //resume listening new messages from the Kafka consumer group 
}; 

//on receipt of message from kafka, push the message to local queue, which then will be processed by worker 
function onMessage(message) { 
    q.push(message, function (err, result) { 
    if (err) { logger.error(err); return }  
    }); 
    consumerGroup.pause(); //Pause kafka consumer group to not receive any more new messages 
} 
関連する問題