· 8 years ago · Dec 23, 2017, 04:04 AM
1const _ = require('lodash');
2const Promise = require('bluebird');
3const moment = require('moment');
4const cassandra = require('cassandra-driver');
5const async = require('async');
6const elasticsearch = require('elasticsearch');
7const axios = require('axios');
8
9let cassandraClient = new cassandra.Client({contactPoints: ['127.0.0.1']});
10let elasticClient = new elasticsearch.Client({host: 'localhost:9200',log: 'trace'});
11cassandraClient.options.socketOptions.readTimeout = 300000;
12
13//ELASTIC SEARCH
14elasticClient.ping({
15 // ping usually has a 3000ms timeout
16 requestTimeout: 1000
17 }, function (error) {
18 if (error) {
19 console.trace('elasticsearch cluster is down!');
20 } else {
21 console.log('All is well');
22 }
23 });
24
25
26
27
28
29let config = {
30 // headers: {"Content-Type": "application/x-ndjson"}
31 headers: {"Content-Type": "application/json"}
32
33};
34
35let ratioCount = 10;
36let insertIntoElastic = (ratio) => {
37 ratioCount++;
38 return elasticClient.index({
39 index: 'event_service',
40 // id: String(ratioCount),
41 type: 'ratios',
42 body: {
43 "Ratio": ratio,
44 "Time": new Date(),
45 }
46 })
47};
48
49//CASSANDRA
50let createKeyspace = () => {
51 return new Promise((resolve,reject) => {
52 cassandraClient.execute("CREATE KEYSPACE events WITH replication = {'class': 'SimpleStrategy', 'replication_factor': '3'};", (err,result) => {
53 if(err){return console.log('ERROR CREATING NAMESPACE')};
54 resolve();
55
56 });
57})
58}
59
60let createTable = () => {
61 return new Promise((resolve,reject) =>{
62 cassandraClient.execute('CREATE TABLE IF NOT EXISTS events.events (id int PRIMARY KEY, listing_id int, time timestamp, type text, metadata text);', (err,result) =>{
63 if (!err){
64 // console.log("EVENT TABLE CREATED OR WAS ALREADY CREATED");
65 } else {
66 console.log("ERROR CREATING EVENT TABLE", err);
67 }
68 resolve();
69 })
70 })
71}
72
73let insertEvent = (id,listing_id,time,type,metadata) => {
74 return new Promise((resolve,reject) => {
75 cassandraClient.execute(`INSERT INTO events.events (id,listing_id,time, type, metadata) VALUES (${id},${listing_id},'${time}','${type}','${metadata}')`, function (err, result) {
76 if (!err){
77 // console.log("EVENT INSERTED INTO THE CASSANDRA DB");
78 resolve();
79 } else {
80 console.log('ERROR INSERTING INTO CASSANDRA DB', err);
81 }
82 // Run next function in series
83 // callback(err, null);
84 })
85 })
86
87}
88
89let describeKeyspaces = () => {
90 return new Promise((resolve,reject) => {
91 cassandraClient.execute('SELECT keyspace_name FROM system_schema.keyspaces;', (err,data) => {
92 if (err) {return console.log('ERROR DESCRIBING KEYSPACE:', err)};
93 // console.log('KEYSPACES DESCRIBED:', data);
94 resolve(data);
95 })
96 // resolve(cassandraClient.metadata.keyspaces);//('events'));
97 })
98}
99
100let seedBatchCassandra = (arr) => {
101 console.log('in the seed batch cassandra function');
102 return cassandraClient.batch(arr, {prepare:true});
103}
104
105let getCassandraCountInt = () => {
106 return cassandraClient.execute("SELECT count(*) FROM events.events WHERE metadata = 'int' ALLOW FILTERING;", [], {autopage:true})
107};
108let getCassandraCountExt = () => {
109 return cassandraClient.execute("SELECT count(*) FROM events.events WHERE metadata = 'ext' ALLOW FILTERING;")
110};
111let esQueue = [];
112let seedBatchElastic = (arr) => {
113 console.log('in the seed batch elastic function');
114
115 esQueue = [];
116 console.log('arr.length in seedBatchElastic', arr.length)
117 arr.forEach((event) => {
118 let curCreate = { "create" : { "_index" : "events2", "_type" : "event_fired2" , "_id":event.id} };
119 let curEvent = { "id" : event.id, "listing_id": event.listing_id, "time": event.time , "type": event.type, "metadata": event.metadata.photo };
120 esQueue.push(curCreate);
121 esQueue.push(curEvent);
122 })
123 console.log('elastic queue', esQueue);
124 console.log('elastic queue length', esQueue.length);
125
126 if (esQueue.length === 1000) {
127 let bulkPost = {body: esQueue};
128 return elasticClient.bulk(bulkPost);
129
130 }
131
132}
133
134let SelectAllEvents = () => {
135 new Promise((resolve,reject) => {
136 resolve(cassandraClient.execute("SELECT * FROM events.events;"));
137 }).then((data) =>{
138 console.log(data);
139 })
140}
141
142module.exports = {
143 // InsertEvent:InsertEvent,
144 // createcassandraClient,createcassandraClient,
145 elasticClient:elasticClient,
146 insertIntoElastic:insertIntoElastic,
147 getCassandraCountExt:getCassandraCountExt,
148 getCassandraCountInt:getCassandraCountInt,
149 seedBatchElastic:seedBatchElastic,
150 seedBatchCassandra:seedBatchCassandra,
151 createKeyspace:createKeyspace,
152 createTable:createTable,
153 insertEvent:insertEvent,
154 SelectAllEvents:SelectAllEvents,
155 describeKeyspaces,describeKeyspaces
156}