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 /index.js | |
| parent | 075c26e94f16379b302570e689d244edef94cb04 (diff) | |
| download | brokentalents-c8d900349a2a518c5fee9610b6a903a2ff546e50.tar.gz brokentalents-c8d900349a2a518c5fee9610b6a903a2ff546e50.zip | |
cache many small periods on disk and aggregate them
Diffstat (limited to 'index.js')
| -rw-r--r-- | index.js | 32 |
1 files changed, 13 insertions, 19 deletions
@@ -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(); |
