-
Notifications
You must be signed in to change notification settings - Fork 0
/
replicator.js
63 lines (49 loc) · 1.43 KB
/
replicator.js
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
// Variables
const SOURCE = process.env.SOURCE || "";
const DESTINATION = process.env.DESTINATION || "";
const TOPIC = process.env.TOPIC || "";
const TIMER = process.env.TIMER || 15000;
const GROUPID = process.env.GROUPID || "";
const ENCODING = process.env.ENCODING || "buffer";
// Modules
var kafka = require('kafka-node');
const utf8 = require('utf8');
const avro = require('avsc');
// Clients
var Consumer = kafka.Consumer;
var Producer = kafka.Producer;
var Client = kafka.KafkaClient;
client_source = new Client({kafkaHost: SOURCE});
client_destination = new Client({kafkaHost: DESTINATION});
// Consumer
consumer = new Consumer(
client_source,
[
{ topic: TOPIC, partition: 0 }
],
{
autoCommit: true,
encoding: ENCODING,
groupId: GROUPID
}
);
// Producer
var producer = new Producer(client_destination, { requireAcks: 1 });
producer.on('ready', function () {
consumer.on('message', function (message) {
producer.send([
{ topic: TOPIC, partition: 0, messages: [message.value], attributes: 0 }
], function (err, result) {
console.log(err || result);
// process.exit();
});
})
consumer.on('error', function (err) {
console.log('error', err);
setTimeout(process.exit(), TIMER);
});
})
producer.on('error', function (err) {
console.log(err);
setTimeout(process.exit(), TIMER);
})