【问题标题】:Kinesis lambda DynamoDBKinesis lambda DynamoDB
【发布时间】:2016-06-24 14:23:53
【问题描述】:

我正在为一个用例学习 AWS 服务。在浏览完文档后,我想出了一个简单的流程。我想使用 Streams API 和 KPL 将数据摄取到 Kinesis 流中。我使用示例 putRecord 方法将数据摄取到流中。我正在将此 JSON 摄取到流中 -

{"userid":1234,"username":"jDoe","firstname":"John","lastname":"Doe"}

一旦数据被摄取,我会在 putRecordResult 中得到以下响应 -

Put Result :{ShardId: shardId-000000000000,SequenceNumber: 49563097246355834103398973318638512162631666140828401666}
Put Result :{ShardId: shardId-000000000000,SequenceNumber: 49563097246355834103398973318645765717549353915876638722}
Put Result :{ShardId: shardId-000000000000,SequenceNumber: 49563097246355834103398973318649392495008197803400757250}

现在我编写了一个 Lambda 函数来获取这些数据并推送到 DynamoDB 表中。这是我的 Lambda 函数 -

console.log('Loading function');
var AWS = require('aws-sdk');
var tableName = "sampleTable";
var doc = require('dynamodb-doc');
var db = new doc.DynamoDB();

exports.handler = (event, context, callback) => {
    //console.log('Received event:', JSON.stringify(event, null, 2));
    event.Records.forEach((record) => {
        // Kinesis data is base64 encoded so decode here
        const payload = new Buffer(record.kinesis.data, 'base64').toString('ascii');
        console.log('Decoded payload:', payload);
        var userid = event.userid;
        var username = event.username;
        var firstname = event.firstname;
        console.log(userid + "," + username +","+ firstname);

        var item = {
            "userid" : userid,
            "username" : username,
            "firstname" : firstname
        };

        var params = {
            TableName : tableName,
            Item : item
        };
        console.log(params);

        db.putItem(params, function(err, data){
            if(err) console.log(err);
            else console.log(data);
        });

    });
    callback(null, `Successfully processed ${event.Records.length} records.`);
};

不知何故,我无法在 lambda 函数执行中看到 console.logs。我在流页面中看到有 putRecord 到流并且得到了,但是不知何故我在 Lambdafunction 页面和 DynamoDB 表中什么都看不到。

我有一个用于将数据摄取到 Kinesis 中的 Java 代码的 IAM 策略,另一个用于 Lambda 函数的 lambda-kinesis-execution-role 以及一个用于 DynamoDB 将数据摄取到表中的策略。

是否有任何教程显示如何以正确的方式完成?我感觉我在这个过程中遗漏了很多要点,例如如何链接所有这些 IAM 策略并使它们同步,以便当数据放入流时它由 Lambda 处理并最终在 Dynamo 中?

非常感谢任何指示和帮助。

【问题讨论】:

  • 您的 Lambda 函数是否被调用?您没有提及,所以我想知道您是否启用了将数据从 Kinesis 传递到您的函数的 AWS Lambda 事件:docs.aws.amazon.com/lambda/latest/dg/…
  • 感谢您的评论。是的,我已经在 lambda 函数的事件源选项卡中添加了 Kinesis 流,它显示启用状态和详细信息为 - 批量大小:100,最后结果:OK。当我使用 Kinesis 示例事件模板和测试配置测试事件时,它给我错误说明项目:{id:未定义,用户名:未定义,名字:未定义}
  • 如果上面的代码是您正在使用的代码的直接副本,那么您引用的是event.userid,但您应该使用payload.userid。您已将 Kinesis 记录解码为 payload 变量。
  • 是的,你是对的。刚刚发现我需要使用 cloudwatch 来查看 console.logs 的输出。我一直在向 lambda 接收数据。只需要在 Cloudwatch 日志中看到这一点 :)

标签: amazon-web-services amazon-dynamodb aws-lambda amazon-kinesis amazon-kinesis-kpl


【解决方案1】:

如果您上面的代码是您正在使用的代码的直接副本,那么您引用的是event.userid,但您应该使用payload.userid。您已将 Kinesis 记录解码为有效负载变量。

【讨论】:

    【解决方案2】:

    您可以使用 Lambda 函数

    1.为 Kinesis 和 Dynamodb 创建 IAM 角色

    2.现在从 dynamodb-process-stream 的蓝图创建一个 Lambda 函数

    3.选择我们从IAM创建的执行角色

    4.点击创建函数 现在转到编辑代码部分并编写以下代码

    const AWS =require('aws-sdk');
    const docClient =new AWS.DynamoDB.DocumentClient({region : 'us-east-1'});
    
    exports.handler = (event, context, callback) => {
       event.Records.forEach((record) => { 
    
        var params={
        Item :{
                ROWTIME:Date.now().toString(),//Dynamodb column name
       DATA:new Buffer(record.kinesis.data, base64').toString('ascii')//Dynamodb column name
                  },
        TableName:'mytable'//Dynamodb Table Name
        };
    
    docClient.put(params,function(err,data){
        if(err){
            callback(err,null);
        }
        else{
            callback(null,data);
        }
    });
    });
    };
    

    【讨论】:

      猜你喜欢
      • 2018-11-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-09-07
      • 2020-10-19
      • 2018-10-25
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多