From 02f703759f0e96a24f029309e49ee7b2edd84e41 Mon Sep 17 00:00:00 2001 From: haydenmcp Date: Fri, 25 Feb 2022 17:31:36 -0500 Subject: [PATCH 1/2] Moved analysis service definition back into index. --- .../src/analysis-client.mjs | 14 +++----------- .../deep-microservice-analysis/src/index.mjs | 9 ++++++++- .../test/analysis-service.test.mjs | 18 ------------------ 3 files changed, 11 insertions(+), 30 deletions(-) diff --git a/packages/deep-microservice-analysis/src/analysis-client.mjs b/packages/deep-microservice-analysis/src/analysis-client.mjs index 5a8cc7a87..254e25648 100644 --- a/packages/deep-microservice-analysis/src/analysis-client.mjs +++ b/packages/deep-microservice-analysis/src/analysis-client.mjs @@ -1,20 +1,15 @@ import {attachExitHandler} from '@thinkdeep/attach-exit-handler'; import {getPublicIP} from '@thinkdeep/get-public-ip'; -import {AnalysisService} from './analysis-service.mjs'; -import {SentimentStore} from './datasource/sentiment-store.mjs'; -import Sentiment from 'sentiment'; class AnalysisClient { - constructor(mongoClient, kafkaClient, apolloServer, expressApp, logger) { + constructor(mongoClient, kafkaClient, apolloServer, expressApp, analysisService, logger) { this._mongoClient = mongoClient; this._kafkaClient = kafkaClient; this._apolloServer = apolloServer; this._expressApp = expressApp; + this._analysisService = analysisService; this._logger = logger; - - this._sentimentStore = undefined; - this._analysisService = undefined; } @@ -37,9 +32,6 @@ class AnalysisClient { await this._mongoClient.connect(); this._logger.info(`Connecting to analysis service.`); - this._sentimentStore = new SentimentStore(this._mongoClient.db('admin').collection('sentiments'), this._logger); - this._analysisService = new AnalysisService(this._sentimentStore, new Sentiment(), this._kafkaClient, this._logger); - await this._analysisService.connect(); } @@ -70,7 +62,7 @@ class AnalysisClient { }); const port = 4001; - await new Promise((resolve) => this._expressApp.listen({port}, resolve)); + await new Promise(((resolve) => this._expressApp.listen({port}, resolve)).bind(this)); this._logger.info(`🚀 Server ready at http://${getPublicIP()}:${port}${this._apolloServer.graphqlPath}`); } } diff --git a/packages/deep-microservice-analysis/src/index.mjs b/packages/deep-microservice-analysis/src/index.mjs index e859bad27..964ed5298 100644 --- a/packages/deep-microservice-analysis/src/index.mjs +++ b/packages/deep-microservice-analysis/src/index.mjs @@ -1,6 +1,8 @@ import {buildSubgraphSchema} from '@apollo/subgraph'; import { AnalysisClient } from './analysis-client.mjs'; +import {AnalysisService} from './analysis-service.mjs'; import {ApolloServer} from 'apollo-server-express'; +import {SentimentStore} from './datasource/sentiment-store.mjs'; import express from 'express'; import { getLogger } from './get-logger.mjs'; import { Kafka } from 'kafkajs'; @@ -8,6 +10,7 @@ import { loggingPlugin } from './logging-plugin.mjs'; import {MongoClient} from 'mongodb'; import {resolvers} from './resolvers.mjs'; import {typeDefs} from './schema.mjs'; +import Sentiment from 'sentiment'; const logger = getLogger(); @@ -20,6 +23,10 @@ const kafkaClient = new Kafka({ brokers: [kafkaBroker] }); + +const sentimentStore = new SentimentStore(mongoClient.db('admin').collection('sentiments'), logger); +const analysisService = new AnalysisService(sentimentStore, new Sentiment(), kafkaClient, logger); + const apolloServer = new ApolloServer({ schema: buildSubgraphSchema([{typeDefs, resolvers}]), dataSources: () => ({analysisService}), @@ -34,7 +41,7 @@ const apolloServer = new ApolloServer({ (async () => { - const microserviceClient = new AnalysisClient(mongoClient, kafkaClient, apolloServer, express(), logger); + const microserviceClient = new AnalysisClient(mongoClient, kafkaClient, apolloServer, express(), analysisService, logger); await microserviceClient.connect(); diff --git a/packages/deep-microservice-analysis/test/analysis-service.test.mjs b/packages/deep-microservice-analysis/test/analysis-service.test.mjs index a1b0b5b8d..5f2837af5 100644 --- a/packages/deep-microservice-analysis/test/analysis-service.test.mjs +++ b/packages/deep-microservice-analysis/test/analysis-service.test.mjs @@ -121,24 +121,6 @@ describe('analysis-service', () => { }) it('should wait for topic creation', async () => { - const economicEntityName = "SomeBusinessName"; - const economicEntityType = "BUSINESS"; - const tweets = [{ - text: 'something' - }, { - text: 'something else' - }]; - const timeSeriesData = [{ - timestamp: 1, - economicEntityName: 'irrelevant', - economicEntityType: 'irrelevant', - tweets - }]; - - const sentimentResult = { - score: 1 - }; - sentimentLib.analyze.returns(sentimentResult); await subject.connect(); From 2bcf79f32afb4a128b7c69be9ce9d76917f79839 Mon Sep 17 00:00:00 2001 From: haydenmcp Date: Fri, 25 Feb 2022 17:33:45 -0500 Subject: [PATCH 2/2] Added fetch interval of 1 minute --- .../deep-microservice-collection/src/collection-service.mjs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/packages/deep-microservice-collection/src/collection-service.mjs b/packages/deep-microservice-collection/src/collection-service.mjs index 87d058734..04bd6f901 100644 --- a/packages/deep-microservice-collection/src/collection-service.mjs +++ b/packages/deep-microservice-collection/src/collection-service.mjs @@ -139,7 +139,9 @@ class CollectionService { * As a result, ~400 businesses can be watched when fetched every 6 hours. */ /** min | hour | day | month | weekday */ - schedule: `0 */6 * * *`, + // schedule: `0 */6 * * *`, + // TODO + schedule: `* * * * *`, image: 'thinkdeeptech/collect-data:latest', command: 'node', args: ['src/collect-data.mjs', `--entity-name=${entityName}`, `--entity-type=${entityType}`, '--operation-type=fetch-tweets']