【问题标题】:Use of Amazon Lambda to inform a user via SNS about a change in the DynamoDB使用 Amazon Lambda 通过 SNS 通知用户有关 DynamoDB 的更改
【发布时间】:2015-12-26 01:56:15
【问题描述】:

目前我正在尝试使用 AWS 作为后台服务构建某种消息传递服务。但是我很挣扎,因为到目前为止我从未使用过 Javascript。

基本思路是:如果DynamoDB有变化,通过推送消息通知“创建者”。

奇怪的是,当我点击“测试”时,它工作得很好,但是当我直接修改数据库时,什么也没有发生。更奇怪的是,一天一两次,当我在数据库中更改它时,我会收到通知。所以这一切似乎有点随机,我希望有人能帮助我。

这是在 lambda 上运行的代码

console.log('Loading function');
var aws = require('aws-sdk');
var doc = require('dynamodb-doc');
var dynamo = new doc.DynamoDB();
var async = require('async');

exports.handler = function(event, context) {


    var target = null;
    var status = null;

    async.waterfall([
            event.Records.forEach(function(record){

                console.log("Function Invoked");

                if(record.eventName == "MODIFY" && record.dynamodb.NewImage.creator.S!=record.dynamodb.NewImage.participant.S){

                    console.log("If-Cause entred!");

                    var creator = record.dynamodb.NewImage.creator.S;
                    var participant = record.dynamodb.NewImage.participant.S;


                    creator = JSON.stringify(creator);
                    status = JSON.stringify(record.dynamodb.NewImage.status.S);
                    creator = creator.replace(/"/g, "");
                    status = status.replace(/"/g, "");

                    participant = JSON.stringify(record.dynamodb.NewImage.participant.S);
                    participant = participant.replace(/"/g, "");



                    console.log("creator: "+ creator);
                    console.log("status: "+ status);

                    console.log("participant: "+ participant);


                    getArn(creator,function(response){
                        console.log("RESPONSE: "+ response);

                        //  var name_part = null;

                        //   getName(participant,function(response){
                        //     name_part = response;
                        //    });

                        var sns = new aws.SNS();
                        var payload_accepted = {
                            "GCM": "{ \"data\": { \"message\": \"Your event has been accepted!!\"} }"
                        };

                        var payload_declined = { "GCM": "{ \"data\":  { \"message\": \"Your event has been declined!\" } }" };

                        //  payload_declined.GCM.data.message = "Hello";



                        payload_accepted = JSON.stringify(payload_accepted);
                        payload_declined = JSON.stringify(payload_declined);

                        var payload;

                        if(status == "true"){
                            payload = payload_accepted;
                        }else{
                            payload = payload_declined
                        }


                        var params = {
                            TargetArn: response,
                            MessageStructure: 'json',
                            Message: payload
                        };

                        sns.publish(
                            params, function(err, data) {
                                if (err) {
                                    console.log(err.stack);

// Notify Lambda that we are finished, but with errors
                                    context.done(err, 'Brians Function Finished with Errors!');
                                    return;
                                }
                                console.log('push sent');
                                console.log(data);

// Notify Lambda that we are finished
                                context.done(null, 'Brians Function Finished!');
                            });



                    });
                }





            })
        ],
        function (err) {

            context.done(null, 'Brians Function Finished!');

            //this last function runs anytime any callback has an error, or if no error
            // then when the last function in the array above invokes callback.
            if (err) { sendForTheCodeDoctor(); }
        });
};


function getArn(userid, callback) {



    var params = {

        TableName : "users",
        Key : {
            "id" : userid
        },

        ProjectionExpression: 'arn'

    }

    dynamo.getItem(params, function(err, data) {


        if (err) {
            console.log(err);
            return err;
        }
        else {
            console.log(data.Item.arn);
            var response = data.Item.arn;

            return callback(response);
        }
    });
}

感谢您的帮助:D

【问题讨论】:

    标签: android node.js amazon-web-services amazon-dynamodb aws-lambda


    【解决方案1】:

    由于您使用async.waterfall 的方式,您有一个竞争条件。你想要的是async.each 或类似的东西(我建议看看the docs 中每个可用函数的作用)。

    目前您的代码将立即同步运行event.Records.forEach 函数,每次迭代都会异步尝试发送一条SNS 消息,但只要forEach 完成(很可能在您的SNS 消息发送之前)将调用async.waterfall 回调并调用context.done 函数。至此,执行结束。即使 SNS 在调用回调之前完成,它本身也会调用 context.done,从而阻止下一次迭代发送 SNS。

    你想要更像这样的东西:

    async.each(event.Records,function(record,done){
      //extract information etc
      sns.publish(params,done);
    }, function(err){
      context.done(err);
    });
    

    这种方式context.done只被调用一次,并且在所有SNS消息都发送完毕后

    【讨论】:

    • 实际上总是必须调用“context.done”。如果没有,你迟早会在 Lambda 上遇到一些奇怪的问题。
    • 确实如此。但是您的代码多次调用它,这是我指出的问题。如果您尝试使用相同的 lambda 调用处理多个事件,则更新后的代码将遇到同样的问题。
    • 请您更正“更新”代码并在此处发布好吗?因为不知何故 Lambda 仍然表现得有些古怪。谢谢!
    【解决方案2】:

    不知何故,我发现一开始的 IF-Condition 造成了麻烦。下面的代码工作得很好,可能会帮助其他人,他们想做和我一样的事情。

    如果 DynamoDB 发生更改,则会自动调用此 Lambda 函数。然后它检查发生了什么样的变化,并通过推送 (SNS) 通知 DynamoDB 条目的创建者。它甚至提供了一个可定制的消息来发送。只需输入你要传输的String即可。

    它基本上是像 Whatsapp 这样的消息服务的核心部分。

    console.log('Loading function');
    var aws = require('aws-sdk');
    var doc = require('dynamodb-doc');
    var dynamo = new doc.DynamoDB();
    
    exports.handler = function(event, context) {    
    
    
        var target = null;
        var status = null;
    
    
        event.Records.forEach(function(record){
    
            console.log("Function Invoked");
            console.log("Function Event: "+record.eventName);
    
    
            if (record.eventName == "MODIFY") db_record_exists(record, function(r) {
              if (r) {
    
                 if(r.handle == "true"){
                    creator = r.creator;
                    participant = r.participant;
    
                    creator = JSON.stringify(creator);
                    status = JSON.stringify(record.dynamodb.NewImage.status.S);
                    creator = creator.replace(/"/g, "");
                    status = status.replace(/"/g, "");
    
                    participant = JSON.stringify(record.dynamodb.NewImage.participant.S);
                    participant = participant.replace(/"/g, "");
    
                    console.log("creator: "+ creator);
                    console.log("status: "+ status);
                    console.log("participant is: "+ participant);    
    
                    getCredentials(creator,participant,function(response){
                    console.log("RESPONSE: "+ response);
    
                    sendSNS(response, status, context);    
    
                    }); 
    
                  }else{
                      shutdown(context);
                  }
              } else shutdown(context);
            }); 
        }
      );
    };
    
    
    
    function db_record_exists(record, callback){
    
            var users = {};
            users.creator = record.dynamodb.NewImage.creator.S;
            users.participant = record.dynamodb.NewImage.participant.S;
    
            if(users.creator != users.participant){
                users.handle = "true";
            }else{
                users.handle = "false";
            }
    
            callback(users);
    } 
    
    function shutdown(context){
            context.done();
    }
    
    
    function sendSNS(creator_information, status, context){
    
            var sns = new aws.SNS();
            var payload_negative = creator_information.name + " has declined your event!";
            var payload_positive = creator_information.name + " has accepted your event!";
           //var payload_positive = "Positive";
    
            var payload_accepted = {
                            "GCM": "{ \"data\": { \"message\": \"Your event has been accepted!!\"} }"
                            };
    
            var payload_declined = { "GCM": "{ \"data\":  { \"message\": \"Your event has been declined!\" } }" };
    
         //   console.log("Standart Message as JSON: "+ payload_declined.GCM);
    
            payload_declined.GCM = "{ \"data\": { \"message\": \""+payload_negative+"\" } }";
            payload_accepted.GCM = "{ \"data\": { \"message\": \""+payload_positive+"\" } }";
    
    
         //   console.log("Standart Message as JSON: "+ payload_declined.GCM);
          //  console.log("Standart Message as String: "+ JSON.stringify(payload_declined));
    
    
            payload_accepted = JSON.stringify(payload_accepted);
            payload_declined = JSON.stringify(payload_declined);
    
            var payload;
    
            if(status == "true"){
                payload = payload_accepted;
            }else{
                payload = payload_declined;
            }
    
    
            var params = {
            TargetArn: creator_information.arn,
            MessageStructure: 'json',
            Message: payload
            };
    
            sns.publish(
               params, function(err, data) {
                if (err) {
                    console.log(err.stack);
    
    // Notify Lambda that we are finished, but with errors
                context.done(err, 'Brians Function Finished with Errors!');  
                return;
            }
            console.log('push sent');
            console.log(data);
    
    // Notify Lambda that we are finished
            context.done(null, 'Brians Function Finished!');  
        });
    } 
    
    function getCredentials(creator,participant,callback) {
    
        var params_name = {
    
        TableName : "users",
        Key : {
            "id" : participant
        },
    
        ProjectionExpression: 'person_name'
    
    }
        var params_arn = {
    
        TableName : "users",
        Key : {
            "id" : creator
        },
    
        ProjectionExpression: 'arn'
    
    }
    
        dynamo.getItem(params_name, function(err, data) {
    
         var response = {};
    
                if (err) {
                    console.log(err);
                    return err;
                }
                else {
                    console.log("creators name is: "+data.Item.person_name);
                     response.name = data.Item.person_name;
                     dynamo.getItem(params_arn, function(err, data) {
    
    
                            if (err) {
                                console.log(err);
                                return err;
                            }
                            else {
                              console.log(data.Item.arn);
    
                           response.arn = data.Item.arn;
    
                                return callback(response);
                            }
                     });
    
    
                }
        });
    
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-10-12
      • 2021-02-22
      • 1970-01-01
      • 1970-01-01
      • 2016-11-30
      • 2019-02-11
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多