Skip to content

Instantly share code, notes, and snippets.

@am-MongoDB
Last active February 13, 2017 16:13
Show Gist options
  • Select an option

  • Save am-MongoDB/0ee9d95fc1abcaf0949245811dcf7d6e to your computer and use it in GitHub Desktop.

Select an option

Save am-MongoDB/0ee9d95fc1abcaf0949245811dcf7d6e to your computer and use it in GitHub Desktop.
Kinesis producer (Node.js)
'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