You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将触发Kinesis流的Lambda函数数据写入Timestream?

将Kinesis触发的Lambda从写入DynamoDB改为写入Timestream的实现方案

现有一个由Kinesis流触发的Lambda函数,原逻辑是将解码后的Kinesis数据写入DynamoDB,现在需要替换为写入Timestream。请问如何用Timestream SDK复现类似的写入流程,或者有没有更简便的实现方式?

原Lambda代码

const AWS = require('aws-sdk');
//const dynamo = new AWS.DynamoDB.DocumentClient();
//const events ="events"

exports.handler = function(event, context) {

    event.Records.forEach(function(record) {
        
        console.log('Data payload: ', record.kinesis.data);
        
        //Kinesis data is base64 encoded so decode here
        var payload = Buffer.from(record.kinesis.data, 'base64').toString('utf8');
        
        console.log('Decoded payload:', payload);
        
        var data = JSON.parse(payload);
        
        console.log("Data: %j", data);
        
        console.log("Data value: %s", data.value);
        
        var clean = JSON.stringify(data.value).replace(/[`~!@#$%^&*()_|+\-=?;'\",.<>{}\[\]\\\/]/gi, '')
        
        console.log("Data clean: %s", clean);
        
        var values = clean.split(":");
        var unit = values[0];
        var metric = values[1];
        
        console.log("Data unit: %s", unit);
        console.log("Data metric: %s", metric);
        
        console.log("Writing records");
        const currentTime = Date.now().toString(); // Unix time in milliseconds
       
        const item = { // 修正原代码语法错误:数组字面量改为对象字面量
            "deviceId": context.awsRequestId,
            "timestamp": data.datetime,
            "time": currentTime.toString(),
            "type": data.topic,
            "metric": metric,
            "forge": "uuid",
            "unit": unit
        };
  
        var params = {
        TableName: events,
        Item:{
            "deviceId": context.awsRequestId,
            "timestamp": data.datetime,
            "type": data.topic,
            "metric": metric,
            "forge": "uuid",
            "unit": unit
            }
        };
        
        
        // Instead of saving data to dynamo. I would like to store it on timestream.
        /*console.log("Saving Telemetry Data");
        
        dynamo.put(params, function(err, data) {
            if (err) {
                console.error("Unable to add event. Error JSON:", JSON.stringify(err, null, 2));
                context.fail();
            } else {
                console.log(data);
                console.log("Data saved:", JSON.stringify(params, null, 2));
                context.succeed();
                return {"message": "Item created in DB"};
            }
        });
        */
    });
};

实现方案

1. 核心改造点

  • 替换DynamoDB客户端为Timestream Write客户端
  • 将原数据结构映射为Timestream的维度+度量时序数据模型
  • 使用批量写入接口提升处理效率
  • 确保Lambda执行角色拥有timestream:WriteRecords权限

2. 修改后的完整代码

const AWS = require('aws-sdk');
// 初始化Timestream Write客户端,替换为你的AWS区域
const timestreamWrite = new AWS.TimestreamWrite({ region: 'your-region' });
// 配置你的Timestream数据库和表名
const DATABASE_NAME = 'your-timestream-db';
const TABLE_NAME = 'your-timestream-table';

exports.handler = async function(event, context) {
    // 批量收集Timestream写入记录
    const records = [];

    for (const record of event.Records) {
        console.log('Data payload: ', record.kinesis.data);
        
        // 解码Kinesis base64数据
        const payload = Buffer.from(record.kinesis.data, 'base64').toString('utf8');
        console.log('Decoded payload:', payload);
        
        const data = JSON.parse(payload);
        console.log("Data: %j", data);
        
        console.log("Data value: %s", data.value);
        const clean = JSON.stringify(data.value).replace(/[`~!@#$%^&*()_|+\-=?;'\",.<>{}\[\]\\\/]/gi, '');
        console.log("Data clean: %s", clean);
        
        const values = clean.split(":");
        const unit = values[0];
        const metric = values[1];
        
        console.log("Data unit: %s", unit);
        console.log("Data metric: %s", metric);

        // 构造Timestream记录:优先使用数据自带的时间戳,无效则用当前时间(Unix毫秒级)
        let timestamp;
        if (data.datetime) {
            timestamp = typeof data.datetime === 'string' 
                ? new Date(data.datetime).getTime().toString() 
                : data.datetime.toString();
        } else {
            timestamp = Date.now().toString();
        }

        const timestreamRecord = {
            // 维度:用于过滤、分组的静态属性
            Dimensions: [
                { Name: 'deviceId', Value: context.awsRequestId },
                { Name: 'type', Value: data.topic },
                { Name: 'unit', Value: unit },
                { Name: 'forge', Value: 'uuid' }
            ],
            // 度量:时序动态数值
            MeasureName: 'metric_value',
            MeasureValue: metric,
            MeasureValueType: 'DOUBLE', // 根据metric类型调整,如VARCHAR、INT等
            Time: timestamp
        };

        records.push(timestreamRecord);
    }

    // 批量写入Timestream
    if (records.length > 0) {
        const params = {
            DatabaseName: DATABASE_NAME,
            TableName: TABLE_NAME,
            Records: records
        };

        try {
            const result = await timestreamWrite.writeRecords(params).promise();
            console.log("数据写入Timestream成功:", JSON.stringify(result, null, 2));
            return { message: `${records.length}条记录写入成功` };
        } catch (err) {
            console.error("写入Timestream失败:", JSON.stringify(err, null, 2));
            throw err;
        }
    }

    return { message: "无有效记录需要写入" };
};

3. 关键说明

  • 数据模型适配:Timestream是时序数据库,需将静态标识属性设为维度,动态数值设为度量,符合时序数据的查询优化逻辑
  • 时间戳要求:Timestream接受Unix毫秒字符串或ISO8601格式时间戳,需确保输入格式合法
  • 权限配置:Lambda执行角色需关联以下IAM策略(替换占位符为你的资源信息):
    {
        "Version": "2012-10-17",
        "Statement": [
            {
                "Effect": "Allow",
                "Action": "timestream:WriteRecords",
                "Resource": "arn:aws:timestream:your-region:your-account-id:database/your-db/table/your-table"
            }
        ]
    }
    
  • 批量写入优势:Kinesis触发Lambda时会传递多条记录,批量写入可减少API调用次数,提升处理效率

内容的提问来源于stack exchange,提问作者fneron

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 15:20:33