Skip to content

Commit 457311b

Browse files
author
Filippo Balicchia
committed
add producer for try plugin
1 parent 3c495f0 commit 457311b

4 files changed

Lines changed: 48 additions & 11 deletions

File tree

bin/producer.js

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
var kafka_native = require('kafka-native');
2+
3+
var broker = 'localhost:9092';
4+
var topic = 'test';
5+
6+
function produceMessage() {
7+
var producer = new kafka_native.Producer({
8+
broker: broker
9+
});
10+
11+
producer.partition_count(topic)
12+
.then(function(npartitions) {
13+
var partition = 0;
14+
setInterval(function() {
15+
for(var i = 0; i < 100; ++i){
16+
message = 'messageProduced-' + i
17+
console.log('producer send message ' + message + ' to partition' + partition);
18+
producer.send(topic, partition, [message]);
19+
partition = (partition + 1) % npartitions;
20+
}
21+
}, 1000);
22+
});
23+
}
24+
25+
produceMessage();

config/examples/kafka-stdout-yml.yml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ input:
77
brokerAddress: localhost
88
offset_directory: /tmp/kafka-offsets
99
brokerPort: 9092
10-
topic: test4
10+
topic: test
1111
debug: true
1212

1313
output:

docker-compose.yml

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,16 @@
1+
zookeeper:
2+
image: wurstmeister/zookeeper
3+
ports:
4+
- 2181:2181
5+
kafka:
6+
image: wurstmeister/kafka:0.9.0.1
7+
ports:
8+
- 9092:9092
9+
links:
10+
- zookeeper
11+
environment:
12+
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
13+
KAFKA_CREATE_TOPICS: "test:1:1"
14+
KAFKA_ADVERTISED_HOST_NAME: 127.0.0.1
15+
KAFKA_ADVERTISED_PORT: 9092
16+
KAFKA_AUTO_CREATE_TOPICS_ENABLE: 'true'

lib/plugins/input/kafka.js

Lines changed: 6 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ var consoleLogger = require('../../util/logger.js')
1111
function InputKafka (config, eventEmitter) {
1212
this.config = config
1313
this.eventEmitter = eventEmitter
14-
this.consumer;
14+
;
1515
}
1616
module.exports = InputKafka
1717
/**
@@ -29,14 +29,11 @@ InputKafka.prototype.start = function () {
2929

3030

3131
InputKafka.prototype.createServer = function () {
32-
33-
consoleLogger.log('start kafka consumer ')
34-
32+
consoleLogger.log('start kafkaConsumer ')
3533
var brokerVar = this.config.brokerAddress + ':' + this.config.brokerPort
3634
var topicVar = this.config.topic;
3735
var eventEmitter = this.eventEmitter
38-
consoleLogger.log('Init kafka consumer ')
39-
consumer = new kafka.Consumer({
36+
global.kafkaConsumer = new kafka.Consumer({
4037
broker: brokerVar,
4138
topic: topicVar,
4239
offset_directory: this.config.offset_directory,
@@ -48,8 +45,7 @@ InputKafka.prototype.createServer = function () {
4845
}
4946
});
5047

51-
consoleLogger.log('start consumer ')
52-
consumer.start()
48+
kafkaConsumer.start()
5349

5450
}
5551

@@ -58,8 +54,8 @@ InputKafka.prototype.createServer = function () {
5854
* we close kafka consumer.
5955
*/
6056
InputKafka.prototype.stop = function (cb) {
61-
consoleLogger.log('kafka input stop')
62-
consumer.stop()
57+
consoleLogger.log('stop kafkaConsumer')
58+
kafkaConsumer.stop()
6359
cb()
6460
}
6561

0 commit comments

Comments
 (0)