如何将AWS IoT Core设备消息传输存储至DocumentDB
AWS IoT Core 消息写入DocumentDB实现方案
IoT Core规则引擎目前确实没有提供DocumentDB的原生直连动作,最通用、生产可用的方案是通过Lambda函数做中转实现,具体步骤如下:
1. 前置环境准备
- 提前部署好DocumentDB集群,建议部署在私有VPC子网内,不要开启公网访问;配置集群安全组入站规则,后续放通Lambda所属安全组的27017(MongoDB协议默认端口)访问权限。
- 提前在DocumentDB中创建好用于存储设备消息的数据库、集合,按需配置好集合的索引,比如按设备ID、上报时间建索引,方便后续查询。
- 准备好Lambda运行所需的IAM角色,需要包含两类权限:一类是VPC访问权限(允许创建、挂载弹性网卡,让Lambda能接入VPC内网访问DocumentDB);另一类是基础的CloudWatch日志写入权限,方便排查问题。
2. 配置IoT Core消息路由规则
- 进入IoT Core控制台,新建消息路由规则:首先编写SQL处理语句,按需筛选、转换需要存储的设备消息,比如要存
devices/+/data主题下的全量上报消息,SQL可以写为SELECT *, topic(2) as deviceId, timestamp() as receiveTime FROM 'devices/+/data',会自动把设备ID、IoT Core接收时间戳附加到消息内容里。 - 给规则添加触发动作,动作类型选择「调用Lambda函数」,绑定你后续部署的处理函数,同时给规则配置对应的IAM权限,允许IoT Core服务调用该Lambda资源。
3. 开发部署Lambda处理函数
- 创建Lambda函数时,选择和DocumentDB相同的VPC,子网选择私有子网,绑定的安全组配置出站规则允许访问DocumentDB的27017端口。
- 函数代码使用MongoDB兼容驱动连接DocumentDB即可,注意一定要做连接复用,不要每次函数触发都新建数据库连接,避免快速耗尽DocumentDB的连接配额。
- 参考代码逻辑(Node.js示例):
const { MongoClient } = require('mongodb'); const DOC_DB_URI = process.env.DOC_DB_URI; // 存在Lambda环境变量里,格式为mongodb://<用户名>:<密码>@<DocumentDB集群内网地址>:27017/ let client; async function initConnection() { if (!client) { client = new MongoClient(DOC_DB_URI, { tls: true, tlsCAFile: `/opt/rds-combined-ca-bundle.pem`, // DocumentDB要求的SSL证书,打包到Lambda层即可 replicaSet: 'rs0', readPreference: 'secondaryPreferred' }); await client.connect(); } return client.db('iot_messages').collection('device_reports'); } exports.handler = async (event) => { try { const collection = await initConnection(); // event就是IoT Core传递过来的处理后的消息内容 await collection.insertOne(event); return { status: 'success' }; } catch (err) { console.error('写入DocumentDB失败', err); throw err; // 抛出错误触发IoT Core的错误重试机制 } }
可选优化方案(针对大消息量场景)
- 如果设备上报消息峰值很高,可以在IoT Core规则和Lambda之间加一层Kinesis数据流做削峰,Lambda按批次消费Kinesis数据后批量写入DocumentDB,降低数据库写入压力。
- 给Lambda配置SQS死信队列,多次重试依然写入失败的消息会自动落入死信队列,后续可以人工补录,避免消息丢失。
- 如果对消息实时性要求不高,也可以选择IoT Core规则先把消息落地到S3,再定时触发批量任务把S3文件批量导入DocumentDB,成本会更低。
内容的提问来源于stack exchange,提问作者Gabriel P
相关产品推荐
相关产品推荐

