From 075c26e94f16379b302570e689d244edef94cb04 Mon Sep 17 00:00:00 2001 From: schneefux Date: Sat, 5 May 2018 18:52:59 +0200 Subject: save data to file --- index.js | 28 ++++++++++++++++++---------- 1 file changed, 18 insertions(+), 10 deletions(-) (limited to 'index.js') diff --git a/index.js b/index.js index cefd7e9..1cf2228 100644 --- a/index.js +++ b/index.js @@ -6,20 +6,28 @@ const R = require('ramda'), source = require('./source'), params = require('./params'), reduce = require('./reduce'), + file = require('./file'), config = require('./config'); -function main() { - const paramsForHourSample = params.paramsForHourSample(config.api); - const loadFApi = source.loadFApi(config.etl); - const aggregatePayloads = reduce.aggregatePayloads(config.reduce); - const deriveStatistics = reduce.deriveStatistics(config.reduce); +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 requestsForHour = (timestamp) => R.map(R.apply(loadFApi), paramsForHourSample(timestamp)); - const requests = R.map(R.apply(loadFApi), paramsForHourSample(moment().subtract(1, 'days'))); +function processFPastHour(config) { + return (timestamp) => + Future.parallel(1, requestsForHour(timestamp)) + .map(R.compose(deriveStatistics, aggregatePayloads, R.unnest)) + .chain(saveFPayloads(timestamp)); +} - /* main */ - Future.parallel(1, requests) - .map(R.compose(deriveStatistics, aggregatePayloads, R.unnest)) - .value(R.compose(console.log, (a) => JSON.stringify(a, null, 2))); +function main() { + const now = moment().subtract(1, 'hours').startOf('hour'); + console.log(requestsForHour(now)); + processFPastHour(config)(now).fork(console.error, console.log); } main(); -- cgit v1.3.1