From c8d900349a2a518c5fee9610b6a903a2ff546e50 Mon Sep 17 00:00:00 2001 From: schneefux Date: Sat, 5 May 2018 23:31:56 +0200 Subject: cache many small periods on disk and aggregate them --- index.js | 32 +++++++++++++------------------- 1 file changed, 13 insertions(+), 19 deletions(-) (limited to 'index.js') 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(); -- cgit v1.3.1