summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2018-05-05 23:31:56 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2018-05-05 23:31:56 +0200
commitc8d900349a2a518c5fee9610b6a903a2ff546e50 (patch)
tree3e99c3f388068e78743653dec46cf106437a7052
parent075c26e94f16379b302570e689d244edef94cb04 (diff)
downloadbrokentalents-c8d900349a2a518c5fee9610b6a903a2ff546e50.tar.gz
brokentalents-c8d900349a2a518c5fee9610b6a903a2ff546e50.zip
cache many small periods on disk and aggregate them
-rw-r--r--config.js33
-rw-r--r--etl.js5
-rw-r--r--file.js13
-rw-r--r--index.js32
-rw-r--r--params.js14
-rw-r--r--reduce.js6
-rw-r--r--source.js38
7 files changed, 105 insertions, 36 deletions
diff --git a/config.js b/config.js
index f2e0b4f..d1601ec 100644
--- a/config.js
+++ b/config.js
@@ -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`,
},
};
diff --git a/etl.js b/etl.js
index bbb8205..7334b7d 100644
--- a/etl.js
+++ b/etl.js
@@ -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,
diff --git a/file.js b/file.js
index 16f2c94..b750921 100644
--- a/file.js
+++ b/file.js
@@ -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,
};
diff --git a/index.js b/index.js
index 1cf2228..70b9f47 100644
--- a/index.js
+++ b/index.js
@@ -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();
diff --git a/params.js b/params.js
index b3dcb98..5edf658 100644
--- a/params.js
+++ b/params.js
@@ -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),
);
diff --git a/reduce.js b/reduce.js
index 8c29055..b17a246 100644
--- a/reduce.js
+++ b/reduce.js
@@ -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,
+};