mirror of
https://github.com/gosticks/DefinitelyTyped.git
synced 2026-10-05 23:37:04 +00:00
added tests for ConsumerGroup and offsets
This commit is contained in:
@@ -138,6 +138,24 @@ hlConsumer.resumeTopics([
|
||||
|
||||
hlConsumer.close(true, function () {});
|
||||
|
||||
var ackBatchOptions = {'noAckBatchSize': 1024, 'noAckBatchAge': 10};
|
||||
var cgOptions: kafka.ConsumerGroupOptions = {
|
||||
host: 'localhost:2181/',
|
||||
batch: ackBatchOptions,
|
||||
groupId: 'groupID',
|
||||
id: 'consumerID',
|
||||
sessionTimeout: 15000,
|
||||
protocol: ["roundrobin"],
|
||||
fromOffset: "latest",
|
||||
migrateHLC: false,
|
||||
migrateRolling: true
|
||||
};
|
||||
|
||||
var consumerGroup = new kafka.ConsumerGroup( cgOptions, ['topic1']);
|
||||
consumerGroup.on('error', (err) => {});
|
||||
consumerGroup.on('message', (msg) => {});
|
||||
consumerGroup.close(true, () => {});
|
||||
|
||||
var offset = new kafka.Offset(basicClient);
|
||||
|
||||
offset.on('ready', function(){});
|
||||
@@ -154,3 +172,6 @@ offset.commit('groupId', [
|
||||
offset.fetchCommits('groupId', [
|
||||
{ topic: 't', partition: 0 }
|
||||
], function (err, data) {});
|
||||
|
||||
offset.fetchLatestOffsets(['t'], (err, offsets) => {})
|
||||
offset.fetchEarliestOffsets(['t'], (err, offsets) => {})
|
||||
|
||||
Reference in New Issue
Block a user