Last active
February 13, 2017 16:13
-
-
Save am-MongoDB/0ee9d95fc1abcaf0949245811dcf7d6e to your computer and use it in GitHub Desktop.
Kinesis producer (Node.js)
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| 'use strict'; | |
| var util = require('util'); | |
| var logger = require('../../util/logger'); | |
| function loggingProducer(kinesis, config) { | |
| var log = logger().getLogger('loggingProducer'); | |
| function _createStreamIfNotCreated(callback) { | |
| var params = { | |
| ShardCount : config.shards, | |
| StreamName : config.stream | |
| }; | |
| kinesis.createStream(params, function(err, data) { | |
| if (err) { | |
| if (err.code !== 'ResourceInUseException') { | |
| callback(err); | |
| return; | |
| } | |
| else { | |
| log.info(util.format('%s stream is already created. Re-using it.', config.stream)); | |
| } | |
| } | |
| else { | |
| log.info(util.format("%s stream doesn't exist. Created a new stream with that name ..", config.stream)); | |
| } | |
| // Poll to make sure stream is in ACTIVE state before start pushing data. | |
| _waitForStreamToBecomeActive(callback); | |
| }); | |
| } | |
| function _waitForStreamToBecomeActive(callback) { | |
| kinesis.describeStream({StreamName : config.stream}, function(err, data) { | |
| if (!err) { | |
| log.info(util.format('Current status of the stream is %s.', data.StreamDescription.StreamStatus)); | |
| if (data.StreamDescription.StreamStatus === 'ACTIVE') { | |
| callback(null); | |
| } | |
| else { | |
| setTimeout(function() { | |
| _waitForStreamToBecomeActive(callback); | |
| }, 1000 * config.waitBetweenDescribeCallsInSeconds); | |
| } | |
| } | |
| }); | |
| } | |
| function _writeToKinesis() { | |
| var currTime = new Date().getMilliseconds(); | |
| var sensor = 'sensor-' + Math.floor(Math.random() * 100000); | |
| var reading = Math.floor(Math.random() * 1000000); | |
| var record = JSON.stringify({ | |
| program: "logging_producer", | |
| time : currTime, | |
| sensor : sensor, | |
| reading : reading | |
| }); | |
| var recordParams = { | |
| Data : record, | |
| PartitionKey : sensor, | |
| StreamName : config.stream | |
| }; | |
| kinesis.putRecord(recordParams, function(err, data) { | |
| if (err) { | |
| log.error(err); | |
| } | |
| else { | |
| log.info('Successfully sent data to Kinesis.'); | |
| } | |
| }); | |
| } | |
| return { | |
| run: function() { | |
| _createStreamIfNotCreated(function(err) { | |
| if (err) { | |
| log.error(util.format('Error creating stream: %s', err)); | |
| return; | |
| } | |
| var count = 0; | |
| while (count < 10) { | |
| setTimeout(_writeToKinesis(), 1000); | |
| count++; | |
| } | |
| }); | |
| } | |
| }; | |
| } | |
| module.exports = loggingProducer; |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment