diff options
| author | schneefux <schneefux+commit@schneefux.xyz> | 2018-05-05 23:31:56 +0200 |
|---|---|---|
| committer | schneefux <schneefux+commit@schneefux.xyz> | 2018-05-05 23:31:56 +0200 |
| commit | c8d900349a2a518c5fee9610b6a903a2ff546e50 (patch) | |
| tree | 3e99c3f388068e78743653dec46cf106437a7052 | |
| parent | 075c26e94f16379b302570e689d244edef94cb04 (diff) | |
| download | brokentalents-c8d900349a2a518c5fee9610b6a903a2ff546e50.tar.gz brokentalents-c8d900349a2a518c5fee9610b6a903a2ff546e50.zip | |
cache many small periods on disk and aggregate them
| -rw-r--r-- | config.js | 33 | ||||
| -rw-r--r-- | etl.js | 5 | ||||
| -rw-r--r-- | file.js | 13 | ||||
| -rw-r--r-- | index.js | 32 | ||||
| -rw-r--r-- | params.js | 14 | ||||
| -rw-r--r-- | reduce.js | 6 | ||||
| -rw-r--r-- | source.js | 38 |
7 files changed, 105 insertions, 36 deletions
@@ -1,23 +1,50 @@ #!/usr/bin/node +const crypto = require('crypto'); + +const sha256 = (x) => crypto.createHash('sha256').update(x, 'utf8').digest('hex'); + +const hashSettings = (config) => sha256( + config.api.regions + + config.api.modes + + config.api.requestsPerInterval + + config.api.interval + + config.api.intervalUnit +).substring(0, 8); + module.exports = { api: { baseUrl: 'https://api.dc01.gamelockerapp.com', - timeout: 5, + timeout: 3, + /* available: eu, na, sg, cn, sa, ea */ regions: ['sg', 'eu', 'na'], - rpm: 5, + /* available: casual, ranked, blitz_pvp_ranked, casual_aral */ modes: ['blitz_pvp_ranked'], + /* available: 2, 3, 4, 5 */ + matchesPerRequest: 5, + /* How many equidistant requests to create + * matchesPerRequest * requestsPerInterval = sample size for that interval */ + requestsPerInterval: 5, + /* the size of a batch */ + interval: 1, + intervalUnit: 'hours', + /* How many batches to request per minute + * requestsPerInterval * intervalsPer Minute <= API key rate limit */ + intervalsPerMinute: 10, }, etl: { + /* attributes to be extracted and renamed from main API data */ participant: [ ['Winner', ['attributes', 'stats', 'winner']] ], roster: [ ], match: [ ], }, reduce: { + /* criteria or 'filters' */ group: ['Actor', 'Talent'], + /* stats */ sum: ['Level', 'Winner'], }, file: { - pattern: (moment) => `./data/${moment.week()}/${moment.toISOString()}.json`, + pattern: (moment) => `./data/${hashSettings(module.exports)}/${moment.toISOString()}.json`, }, }; @@ -157,7 +157,10 @@ function loadFPayloads(config) { Future.both, R.map(R.pipe( joinParticipantRosterMatchPayloadsToTelemetryPayloads, - R.map(R.omit(['_asset_id', '_participant_id', '_participant_ids', '_roster_id', '_roster_ids', '_match_id', 'Team'])), + R.map(R.pipe( + R.omit(['_asset_id', '_participant_id', '_participant_ids', '_roster_id', '_roster_ids', '_match_id', 'Team']), + R.set(R.lensProp('Count'), 1), + )), )), ), [ loadFParticipantRosterMatchPayloads, @@ -1,12 +1,13 @@ #!/usr/bin/node const R = require('ramda'), - Future = require('fluture') + Future = require('fluture'), + fs = require('fs'), fsPath = require('fs-path'); function loadFPayloads(config) { return R.curry((file) => - Future.node((cb) => fsPath.readFile(file, 'utf8', cb)) + Future.node((cb) => fs.readFile(file, 'utf8', cb)) .chain(Future.encase(JSON.parse)) ); } @@ -14,6 +15,7 @@ function loadFPayloads(config) { function saveFPayloads(config) { return R.curry((file, data) => Future.node((cb) => fsPath.writeFile(file, JSON.stringify(data), cb)) + .map(() => data) ); } @@ -23,8 +25,15 @@ function saveFPayloadsTimestamped(config) { ); } +function loadFPayloadsTimestamped(config) { + return R.curry((timestamp) => + loadFPayloads(config)(config.pattern(timestamp)) + ); +} + module.exports = { loadFPayloads, + loadFPayloadsTimestamped, saveFPayloads, saveFPayloadsTimestamped, }; @@ -3,31 +3,25 @@ const R = require('ramda'), Future = require('fluture'), moment = require('moment'), - source = require('./source'), - params = require('./params'), reduce = require('./reduce'), - file = require('./file'), + source = require('./source'), config = require('./config'); -const paramsForHourSample = params.paramsForHourSample(config.api); -const loadFApi = source.loadFApi(config.etl); -const saveFPayloads = file.saveFPayloadsTimestamped(config.file); -const aggregatePayloads = reduce.aggregatePayloads(config.reduce); -const deriveStatistics = reduce.deriveStatistics(config.reduce); +const loadFTimestamped = source.loadFTimestamped(config); +const deriveStatistics = reduce.deriveStatistics(config.reduce); +const aggregatePayloads = reduce.aggregatePayloads(config.reduce); -const requestsForHour = (timestamp) => R.map(R.apply(loadFApi), paramsForHourSample(timestamp)); +function main() { + const firstOfMay = moment('2018-05-01'); -function processFPastHour(config) { - return (timestamp) => - Future.parallel(1, requestsForHour(timestamp)) - .map(R.compose(deriveStatistics, aggregatePayloads, R.unnest)) - .chain(saveFPayloads(timestamp)); -} + const later = R.curry((base, hs) => base.clone().add(hs, 'hours')); + const hours = R.range(0, 24 * 4); + const laterMoments = R.map(later(firstOfMay), hours); -function main() { - const now = moment().subtract(1, 'hours').startOf('hour'); - console.log(requestsForHour(now)); - processFPastHour(config)(now).fork(console.error, console.log); + const futures = R.map(loadFTimestamped, laterMoments); + Future.parallel(1, futures) + .map(R.compose(deriveStatistics, aggregatePayloads, R.unnest)) + .fork(console.error, (d) => console.log(JSON.stringify(d, null, 2))); } main(); @@ -1,6 +1,7 @@ #!/usr/bin/node -const R = require('ramda'); +const R = require('ramda'), + moment = require('moment'); function paramsForHourSample(config) { const requestArgsForRegionStartEnd = R.curry((region, start, end) => [ @@ -9,17 +10,18 @@ function paramsForHourSample(config) { 'filter[gameMode]': R.join(',', config.modes), 'filter[createdAt-start]': start.toISOString(), 'filter[createdAt-end]': end.toISOString(), - 'page[limit]': '5', + 'page[limit]': config.matchesPerRequest, 'page[offset]': '0', } ]); - const requestsPerMinutePerRegion = Math.floor(config.rpm / config.regions.length); - const splitDurationMinutes = Math.floor(60 / requestsPerMinutePerRegion); + const intervalMinutes = moment.duration(config.interval, config.intervalUnit).as('minutes'); + const requestsPerMinutePerRegion = Math.floor(config.requestsPerInterval / config.regions.length); + const splitDurationMinutes = Math.floor(intervalMinutes / requestsPerMinutePerRegion); const splitMoments = (hour) => R.map( (offsetIndex) => [ - hour.clone().startOf('hour').minutes(offsetIndex * splitDurationMinutes), - hour.clone().startOf('hour').minutes((offsetIndex + 1) * splitDurationMinutes), + hour.clone().minutes(offsetIndex * splitDurationMinutes), + hour.clone().minutes((offsetIndex + 1) * splitDurationMinutes), ], R.range(0, requestsPerMinutePerRegion), ); @@ -19,14 +19,10 @@ function aggregatePayloads(config) { R.prop(0), ), R.pipe( - R.project(config.sum), + R.project(R.concat(config.sum, R.of('Count'))), utils.mergeAllWith(R.add), utils.unlist, ), - R.pipe( - R.length, - R.objOf('Count'), - ), ]), ), ); diff --git a/source.js b/source.js new file mode 100644 index 0000000..2d02939 --- /dev/null +++ b/source.js @@ -0,0 +1,38 @@ +#!/usr/bin/node + +const R = require('ramda'), + Future = require('fluture'), + etl = require('./etl'), + params = require('./params'), + reduce = require('./reduce'), + file = require('./file'); + +function doFetchProcessStoreHour(config) { + const intervalSleep = 60 / config.api.intervalsPerMinute * 1000; + + const paramsForHourSample = params.paramsForHourSample(config.api); + const loadFApi = etl.loadFPayloads(config.etl); + const saveFPayloads = file.saveFPayloadsTimestamped(config.file); + const aggregatePayloads = reduce.aggregatePayloads(config.reduce); + + const requestsForHour = (timestamp) => R.map(R.apply(loadFApi), paramsForHourSample(timestamp)); + + return (timestamp) => Future.parallel(1, requestsForHour(timestamp)) + .map(R.compose(aggregatePayloads, R.unnest)) + .chain(/*R.map(R.chain(Future.after(intervalSleep),*/ saveFPayloads(timestamp))/*))*/ + .chain(Future.after(intervalSleep)); + // TODO sleeps after request; throttling rpm should happen + // at a higher level though. CPU / IO / AWS requests stall + // as well +} + +function loadFTimestamped(config) { + // load from file or fall back to API and store + return (timestamp) => + file.loadFPayloads(config)(config.file.pattern(timestamp)) + .or(doFetchProcessStoreHour(config)(timestamp)); +}; + +module.exports = { + loadFTimestamped, +}; |
