如何将触发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
相关产品推荐
相关产品推荐

