From 73bdb340856e3c4bb50a8b379b38074989159b86 Mon Sep 17 00:00:00 2001 From: CityofDenver <34748483+CityofDenver@users.noreply.github.com> Date: Sun, 7 Jan 2018 21:56:58 +0530 Subject: [PATCH] Lambda function to process Waze Data from S3 This function is triggered by S3 when a new Waze data is populated from waze-downloader lambda function. --- code/lambda-functions/package.json | 7 ++++ code/lambda-functions/waze-processor.js | 55 +++++++++++++++++++++++++ 2 files changed, 62 insertions(+) create mode 100644 code/lambda-functions/package.json create mode 100644 code/lambda-functions/waze-processor.js diff --git a/code/lambda-functions/package.json b/code/lambda-functions/package.json new file mode 100644 index 0000000..dd1c343 --- /dev/null +++ b/code/lambda-functions/package.json @@ -0,0 +1,7 @@ +{ + "name": "waze-processor", + "version": "0.0.1", + "dependencies": { + "aws-sdk": "^2.177.0" + } +} diff --git a/code/lambda-functions/waze-processor.js b/code/lambda-functions/waze-processor.js new file mode 100644 index 0000000..4e2c984 --- /dev/null +++ b/code/lambda-functions/waze-processor.js @@ -0,0 +1,55 @@ +const AWS = require('aws-sdk'); +const s3 = new AWS.S3(); +const sns = new AWS.SNS(); + +exports.handler = (event, context, callback) => { + console.log('Received event:', JSON.stringify(event)); + const bucket = event.Records[0].s3.bucket.name; + const key = decodeURIComponent(event.Records[0].s3.object.key.replace(/\+/g, ' ')); + const params = { + Bucket: bucket, + Key: key, + }; + console.log('Bucket: ' + bucket); + console.log('Key: ' + key); + s3.getObject(params, (err, data) => { + if (err) { + const message = `Error getting object ${key} from bucket ${bucket}. Reason:` + err; + console.log(message); + callback(message); + } else { + let wazeDataObj = data.Body.toString('utf-8'); + const pWazeDataObj = JSON.parse(wazeDataObj); + const alerts = pWazeDataObj.alerts; + const jams = pWazeDataObj.jams; + if (wazeDataObj !== null && wazeDataObj !== undefined) { + console.log(alerts); + const alertsParams = { + Message: JSON.stringify(alerts), + TopicArn: process.env.ALERTSTOPIC + }; + sns.publish(alertsParams, function (err, data) { + if (err) { + console.log(err, err.stack); + } + else { + console.log('--------------------'); + console.log(data); + if (jams !== undefined && jams !== null) { + const jamParams = { + Message: JSON.stringify(jams), + TopicArn: process.env.JAMSTOPIC + }; + sns.publish(jamParams, function (err, data) { + if (err) console.log(err, err.stack); + else console.log(data); + }); + + } + } + }); + } + } + }); +}; +