From 3d0ce078e1edb39b2abbcd9ecfbcb81d6bf4f055 Mon Sep 17 00:00:00 2001 From: Sahil Palvia Date: Thu, 20 Jul 2017 15:08:12 -0700 Subject: [PATCH 1/5] Updating libraries and including shutdownRequested feature. --- bin/kcl-bootstrap | 14 ++++----- lib/kcl/kcl_manager.js | 29 ++++++++++++++++++- lib/kcl/kcl_process.js | 3 +- .../basic_sample/consumer/sample_kcl_app.js | 7 +++++ 4 files changed, 44 insertions(+), 9 deletions(-) diff --git a/bin/kcl-bootstrap b/bin/kcl-bootstrap index 4c24ed32..9643bb04 100755 --- a/bin/kcl-bootstrap +++ b/bin/kcl-bootstrap @@ -33,13 +33,13 @@ var MAVEN_PACKAGE_LIST = [ getMavenPackageInfo('commons-logging', 'commons-logging', '1.1.3'), getMavenPackageInfo('commons-lang', 'commons-lang', '2.6'), getMavenPackageInfo('joda-time', 'joda-time', '2.8.1'), - getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-core', '1.11.14'), - getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-cloudwatch', '1.11.14'), - getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-dynamodb', '1.11.14'), - getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kinesis', '1.11.14'), - getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kms', '1.11.14'), - getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-s3', '1.11.14'), - getMavenPackageInfo('com.amazonaws', 'amazon-kinesis-client', '1.7.2'), + getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-core', '1.11.151'), + getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-cloudwatch', '1.11.151'), + getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-dynamodb', '1.11.151'), + getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kinesis', '1.11.151'), + getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kms', '1.11.151'), + getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-s3', '1.11.151'), + getMavenPackageInfo('com.amazonaws', 'amazon-kinesis-client', '1.7.6'), getMavenPackageInfo('com.fasterxml.jackson.core', 'jackson-databind', '2.6.6'), getMavenPackageInfo('com.fasterxml.jackson.core', 'jackson-core', '2.6.6'), getMavenPackageInfo('com.fasterxml.jackson.core', 'jackson-annotations', '2.6.0'), diff --git a/lib/kcl/kcl_manager.js b/lib/kcl/kcl_manager.js index 2f4a107b..44f564c3 100644 --- a/lib/kcl/kcl_manager.js +++ b/lib/kcl/kcl_manager.js @@ -74,6 +74,10 @@ var KCLStateMachine = BehavioralFsm.extend({ this.transition(context, 'Processing'); return true; }, + beginShutdownRequested: function(context) { + this.transition(context, 'ShutdownRequested'); + return true; + }, beginShutdown: function(context) { this.transition(context, 'ShuttingDown'); return true; @@ -107,6 +111,20 @@ var KCLStateMachine = BehavioralFsm.extend({ return false; } }, + ShutdownRequested: { + beginCheckpoint: function(context) { + this.transition(context, 'Checkpointing'); + return true; + }, + finishShutdownRequested: function(context) { + this.transition(context, 'Ready'); + return true; + }, + '*': function(context) { + this.transition(context, 'Error'); + return false; + } + }, ShuttingDown: { beginCheckpoint: function(context) { this.transition(context, 'FinalCheckpointing'); @@ -212,7 +230,10 @@ KCLManager.prototype.checkpoint = function(sequenceNumber) { */ KCLManager.prototype._onAction = function(action) { var actionType = action.action; - if (actionType === 'initialize' || actionType === 'processRecords' || actionType === 'shutdown') { + if (actionType === 'initialize' || + actionType === 'processRecords' || + actionType === 'shutdown' || + actionType === 'shutdownRequested') { this._onRecordProcessorAction(action); } else if (actionType === 'checkpoint') { @@ -265,6 +286,12 @@ KCLManager.prototype._onRecordProcessorAction = function(action) { beginActionInput = 'beginShutdown'; finishActionInput = 'finishShutdown'; } + else if (actionType === 'shutdownRequested') { + recordProcessorFuncInput.checkpointer = checkpointer; + recordProcessorFunc = recordProcessor.shutdownRequested; + beginActionInput = 'beginShutdownRequested'; + finishActionInput = 'finishShutdownRequested'; + } // Should not occur. else { throw new Error(util.format('Invalid action for record processor: %j', action)); diff --git a/lib/kcl/kcl_process.js b/lib/kcl/kcl_process.js index 495b67f4..6b8b4730 100644 --- a/lib/kcl/kcl_process.js +++ b/lib/kcl/kcl_process.js @@ -69,7 +69,8 @@ var KCLManager = require('./kcl_manager'); function KCLProcess(recordProcessor, inputFile, outputFile, errorFile) { if (typeof recordProcessor.initialize !== 'function' || typeof recordProcessor.processRecords !== 'function' || - typeof recordProcessor.shutdown !== 'function') { + typeof recordProcessor.shutdown !== 'function' || + typeof recordProcessor.shutdownRequested !== 'function') { throw new Error('Record processor must implement initialize, processRecords, and shutdown functions.'); } inputFile = typeof inputFile !== 'undefined' ? inputFile : process.stdin; diff --git a/samples/basic_sample/consumer/sample_kcl_app.js b/samples/basic_sample/consumer/sample_kcl_app.js index 8afd8ff6..aad62aba 100644 --- a/samples/basic_sample/consumer/sample_kcl_app.js +++ b/samples/basic_sample/consumer/sample_kcl_app.js @@ -66,6 +66,13 @@ function recordProcessor() { }); }, + shutdownRequested: function(shutdownRequestedInput, completeCallback) { + log.info('Shutdown requested called.') + shutdownRequestedInput.checkpointer.checkpoint(function (err) { + completeCallback(); + }); + }, + shutdown: function(shutdownInput, completeCallback) { // Checkpoint should only be performed when shutdown reason is TERMINATE. if (shutdownInput.reason !== 'TERMINATE') { From 250461ba09ca5e895037e9f18fa612dd0b29e054 Mon Sep 17 00:00:00 2001 From: Sahil Palvia Date: Thu, 20 Jul 2017 16:08:29 -0700 Subject: [PATCH 2/5] Adding logging statements. --- samples/basic_sample/consumer/sample_kcl_app.js | 1 + 1 file changed, 1 insertion(+) diff --git a/samples/basic_sample/consumer/sample_kcl_app.js b/samples/basic_sample/consumer/sample_kcl_app.js index aad62aba..aca5fef5 100644 --- a/samples/basic_sample/consumer/sample_kcl_app.js +++ b/samples/basic_sample/consumer/sample_kcl_app.js @@ -74,6 +74,7 @@ function recordProcessor() { }, shutdown: function(shutdownInput, completeCallback) { + log.info('Shutdown is called.') // Checkpoint should only be performed when shutdown reason is TERMINATE. if (shutdownInput.reason !== 'TERMINATE') { completeCallback(); From de7ef5d5a4095ddbd9c6f05d53ab5a4c8036e76c Mon Sep 17 00:00:00 2001 From: Sahil Palvia Date: Mon, 24 Jul 2017 16:07:57 -0700 Subject: [PATCH 3/5] Fixing shutdownRequested feature. --- bin/kcl-bootstrap | 8 ++- lib/kcl/kcl_manager.js | 51 +++++++++++++++---- lib/kcl/kcl_process.js | 3 +- .../basic_sample/consumer/sample.properties | 2 +- .../basic_sample/consumer/sample_kcl_app.js | 2 - 5 files changed, 51 insertions(+), 15 deletions(-) diff --git a/bin/kcl-bootstrap b/bin/kcl-bootstrap index 9643bb04..81a642e2 100755 --- a/bin/kcl-bootstrap +++ b/bin/kcl-bootstrap @@ -39,6 +39,7 @@ var MAVEN_PACKAGE_LIST = [ getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kinesis', '1.11.151'), getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kms', '1.11.151'), getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-s3', '1.11.151'), + getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-sts', '1.11.151'), getMavenPackageInfo('com.amazonaws', 'amazon-kinesis-client', '1.7.6'), getMavenPackageInfo('com.fasterxml.jackson.core', 'jackson-databind', '2.6.6'), getMavenPackageInfo('com.fasterxml.jackson.core', 'jackson-core', '2.6.6'), @@ -47,7 +48,12 @@ var MAVEN_PACKAGE_LIST = [ getMavenPackageInfo('org.apache.httpcomponents', 'httpclient', '4.5.2'), getMavenPackageInfo('org.apache.httpcomponents', 'httpcore', '4.4.4'), getMavenPackageInfo('com.google.guava', 'guava', '18.0'), - getMavenPackageInfo('com.google.protobuf', 'protobuf-java', '2.6.1') + getMavenPackageInfo('com.google.protobuf', 'protobuf-java', '2.6.1'), + getMavenPackageInfo('org.slf4j', 'log4j-over-slf4j', '1.7.21'), + getMavenPackageInfo('org.slf4j', 'slf4j-simple', '1.7.21'), + getMavenPackageInfo('org.slf4j', 'slf4j-api', '1.7.21'), + getMavenPackageInfo('org.slf4j', 'jcl-over-slf4j', '1.7.21') + // getMavenPackageInfo('ch.qos.logback', 'logback-classic', '1.1.7') ]; var DEFAULT_JAR_PATH = path.resolve(path.join(__dirname, '..', 'lib', 'jars')); diff --git a/lib/kcl/kcl_manager.js b/lib/kcl/kcl_manager.js index 44f564c3..8f1753c0 100644 --- a/lib/kcl/kcl_manager.js +++ b/lib/kcl/kcl_manager.js @@ -111,9 +111,19 @@ var KCLStateMachine = BehavioralFsm.extend({ return false; } }, + ShutdownRequestedCheckpointing: { + finishCheckpoint: function(context) { + this.transition(context, 'ShutdownRequested'); + return true; + }, + '*': function(context) { + this.transition(context, 'Error'); + return false; + } + }, ShutdownRequested: { beginCheckpoint: function(context) { - this.transition(context, 'Checkpointing'); + this.transition(context, 'ShutdownRequestedCheckpointing'); return true; }, finishShutdownRequested: function(context) { @@ -232,13 +242,15 @@ KCLManager.prototype._onAction = function(action) { var actionType = action.action; if (actionType === 'initialize' || actionType === 'processRecords' || - actionType === 'shutdown' || - actionType === 'shutdownRequested') { + actionType === 'shutdown') { this._onRecordProcessorAction(action); } else if (actionType === 'checkpoint') { this._onCheckpointAction(action); } + else if (actionType === 'shutdownRequested') { + this._onShutdownRequested(action); + } else { this._reportError(util.format('Invalid action received: %j', action)); } @@ -286,12 +298,6 @@ KCLManager.prototype._onRecordProcessorAction = function(action) { beginActionInput = 'beginShutdown'; finishActionInput = 'finishShutdown'; } - else if (actionType === 'shutdownRequested') { - recordProcessorFuncInput.checkpointer = checkpointer; - recordProcessorFunc = recordProcessor.shutdownRequested; - beginActionInput = 'beginShutdownRequested'; - finishActionInput = 'finishShutdownRequested'; - } // Should not occur. else { throw new Error(util.format('Invalid action for record processor: %j', action)); @@ -336,6 +342,33 @@ KCLManager.prototype._onCheckpointAction = function(action) { checkpointer.onCheckpointerResponse.apply(checkpointer, [action.error, action.sequenceNumber]); }; +/** + * Gets invoked when shutdownRequested is called. + * @param {Object} action - RecordProcessor related action + * @private + */ +KCLManager.prototype._onShutdownRequested = function(action) { + var context = this._context; + var recordProcessor = context.recordProcessor; + var recordProcessorFunc = recordProcessor.shutdownRequested; + + if (typeof recordProcessorFunc === 'function') { + var recordProcessorFuncInput = cloneToInput(action); + var checkpointer = context.checkpointer; + + this._handleStateInput(context, 'beginShutdownRequested'); + var callbackFunc = function() { + this._recordProcessorCallback(context, action, 'finishShutdownRequested'); + }.bind(this); + + recordProcessorFuncInput.checkpointer = checkpointer; + recordProcessorFunc.apply(recordProcessor, [recordProcessorFuncInput, callbackFunc]); + } + else { + this._sendAction(context, {action: 'status', responseFor: action.action}); + } +}; + /** * Sends the given action to the MultiLangDaemon. * @param {object} context - Record processor context for which this action belongs to. diff --git a/lib/kcl/kcl_process.js b/lib/kcl/kcl_process.js index 6b8b4730..495b67f4 100644 --- a/lib/kcl/kcl_process.js +++ b/lib/kcl/kcl_process.js @@ -69,8 +69,7 @@ var KCLManager = require('./kcl_manager'); function KCLProcess(recordProcessor, inputFile, outputFile, errorFile) { if (typeof recordProcessor.initialize !== 'function' || typeof recordProcessor.processRecords !== 'function' || - typeof recordProcessor.shutdown !== 'function' || - typeof recordProcessor.shutdownRequested !== 'function') { + typeof recordProcessor.shutdown !== 'function') { throw new Error('Record processor must implement initialize, processRecords, and shutdown functions.'); } inputFile = typeof inputFile !== 'undefined' ? inputFile : process.stdin; diff --git a/samples/basic_sample/consumer/sample.properties b/samples/basic_sample/consumer/sample.properties index 00a1c8a3..f5f10c91 100644 --- a/samples/basic_sample/consumer/sample.properties +++ b/samples/basic_sample/consumer/sample.properties @@ -1,7 +1,7 @@ # The script that abides by the multi-language protocol. This script will # be executed by the MultiLangDaemon, which will communicate with this script # over STDIN and STDOUT according to the multi-language protocol. -executableName = node sample_kcl_app.js +executableName = sample_kcl_app.js # The name of an Amazon Kinesis stream to process. streamName = kclnodejssample diff --git a/samples/basic_sample/consumer/sample_kcl_app.js b/samples/basic_sample/consumer/sample_kcl_app.js index aca5fef5..4c3d89db 100644 --- a/samples/basic_sample/consumer/sample_kcl_app.js +++ b/samples/basic_sample/consumer/sample_kcl_app.js @@ -67,14 +67,12 @@ function recordProcessor() { }, shutdownRequested: function(shutdownRequestedInput, completeCallback) { - log.info('Shutdown requested called.') shutdownRequestedInput.checkpointer.checkpoint(function (err) { completeCallback(); }); }, shutdown: function(shutdownInput, completeCallback) { - log.info('Shutdown is called.') // Checkpoint should only be performed when shutdown reason is TERMINATE. if (shutdownInput.reason !== 'TERMINATE') { completeCallback(); From 0aae883ca1fced0bf82730698c80aa5025cf422e Mon Sep 17 00:00:00 2001 From: Sahil Palvia Date: Mon, 24 Jul 2017 16:11:01 -0700 Subject: [PATCH 4/5] Updating the required libraries list. --- bin/kcl-bootstrap | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/bin/kcl-bootstrap b/bin/kcl-bootstrap index 81a642e2..9643bb04 100755 --- a/bin/kcl-bootstrap +++ b/bin/kcl-bootstrap @@ -39,7 +39,6 @@ var MAVEN_PACKAGE_LIST = [ getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kinesis', '1.11.151'), getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-kms', '1.11.151'), getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-s3', '1.11.151'), - getMavenPackageInfo('com.amazonaws', 'aws-java-sdk-sts', '1.11.151'), getMavenPackageInfo('com.amazonaws', 'amazon-kinesis-client', '1.7.6'), getMavenPackageInfo('com.fasterxml.jackson.core', 'jackson-databind', '2.6.6'), getMavenPackageInfo('com.fasterxml.jackson.core', 'jackson-core', '2.6.6'), @@ -48,12 +47,7 @@ var MAVEN_PACKAGE_LIST = [ getMavenPackageInfo('org.apache.httpcomponents', 'httpclient', '4.5.2'), getMavenPackageInfo('org.apache.httpcomponents', 'httpcore', '4.4.4'), getMavenPackageInfo('com.google.guava', 'guava', '18.0'), - getMavenPackageInfo('com.google.protobuf', 'protobuf-java', '2.6.1'), - getMavenPackageInfo('org.slf4j', 'log4j-over-slf4j', '1.7.21'), - getMavenPackageInfo('org.slf4j', 'slf4j-simple', '1.7.21'), - getMavenPackageInfo('org.slf4j', 'slf4j-api', '1.7.21'), - getMavenPackageInfo('org.slf4j', 'jcl-over-slf4j', '1.7.21') - // getMavenPackageInfo('ch.qos.logback', 'logback-classic', '1.1.7') + getMavenPackageInfo('com.google.protobuf', 'protobuf-java', '2.6.1') ]; var DEFAULT_JAR_PATH = path.resolve(path.join(__dirname, '..', 'lib', 'jars')); From 60ef05bfa5671680bfb6cea50c781a4c200ede63 Mon Sep 17 00:00:00 2001 From: Sahil Palvia Date: Mon, 24 Jul 2017 16:12:52 -0700 Subject: [PATCH 5/5] Updating the executable name in the properties file. --- samples/basic_sample/consumer/sample.properties | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/samples/basic_sample/consumer/sample.properties b/samples/basic_sample/consumer/sample.properties index f5f10c91..00a1c8a3 100644 --- a/samples/basic_sample/consumer/sample.properties +++ b/samples/basic_sample/consumer/sample.properties @@ -1,7 +1,7 @@ # The script that abides by the multi-language protocol. This script will # be executed by the MultiLangDaemon, which will communicate with this script # over STDIN and STDOUT according to the multi-language protocol. -executableName = sample_kcl_app.js +executableName = node sample_kcl_app.js # The name of an Amazon Kinesis stream to process. streamName = kclnodejssample