-
Notifications
You must be signed in to change notification settings - Fork 1
/
listener.js
59 lines (50 loc) · 1.47 KB
/
listener.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
var amqp = require('amqplib')
// const AMQP_URL = ''
const queueName = 'test'
const cuid = require('cuid');
const elasticSearchTemplate = require('./elasticsearch-template.json')
var elasticsearch = require('elasticsearch');
var client = new elasticsearch.Client({
host: 'localhost:9200'
});
client.indices.putTemplate({
name: 'mytest',
body: elasticSearchTemplate
}, (err, result) => {
if (err) console.log(err)
// if (result) console.log(result)
})
amqp.connect(AMQP_URL).then(async(conn) => {
console.log('amqp.connected!')
const channel = await conn.createChannel()
await channel.assertQueue(queueName)
await channel.prefetch(2)
let count = 0;
channel.consume(queueName, async(msg) => {
const json = JSON.parse(msg.content.toString())
count += 1;
try {
var result = await createDocument(json)
console.log(`${count} - ${result}`)
} catch (err) {
// console.log(result)
console.log(err)
}
channel.ack(msg);
})
}).catch((err) => {
console.log(err)
});
async function createDocument(json) {
return new Promise((resolve, reject) => {
client.create({
index: "mytest",
type: 'mytype',
id : cuid(),
body : json
}, function (err, response) {
if (err) return reject(err)
if (response) return resolve(response.result + ' ' + response._id)
})
});
}