如何用Python实现Lambda从AWS IoT Core取数并写入RDS MySQL
将AWS IoT Core数据写入MySQL的实现方案
方案一:AWS Lambda + IoT Core 规则(近实时推荐)
前置准备
- 确保MySQL数据库(AWS RDS或自建)可被Lambda访问:
- 若用AWS RDS,将Lambda配置到RDS所在VPC,同时在RDS安全组中允许Lambda所在安全组的3306端口入站流量。
- 若为自建MySQL,需开放3306端口给Lambda的IP段,或使用公网地址(注意限制IP范围保障安全)。
- 在AWS Secrets Manager存储MySQL账号、密码、地址、端口信息,避免硬编码(Lambda需配置Secrets Manager读取权限)。
- 创建Lambda执行角色,添加以下权限:
AWSIoTFullAccess(或更细粒度的IoT规则触发权限)SecretsManagerRead(仅需读取权限)AmazonRDSDataFullAccess(使用RDS时配置)
Lambda Python 示例代码
import pymysql import json import boto3 def get_mysql_config(secret_name): # 从Secrets Manager获取MySQL配置 client = boto3.client('secretsmanager') response = client.get_secret_value(SecretId=secret_name) secret = json.loads(response['SecretString']) return { 'host': secret['host'], 'user': secret['username'], 'password': secret['password'], 'database': secret['dbname'], 'port': int(secret['port']) } def lambda_handler(event, context): # 解析IoT Core传来的MQTT消息 payload = json.loads(event['records'][0]['value']) # 替换为你的Secrets Manager密钥名称 mysql_config = get_mysql_config('your-mysql-secret-name') try: # 连接MySQL conn = pymysql.connect(**mysql_config) cursor = conn.cursor() # 替换为你的表结构与字段,示例表为iot_data,含timestamp、device_id、temperature字段 insert_sql = """ INSERT INTO iot_data (timestamp, device_id, temperature) VALUES (%s, %s, %s) """ # 从payload提取对应字段,需匹配你的实际消息结构 data = (payload['timestamp'], payload['device_id'], payload['temperature']) cursor.execute(insert_sql, data) conn.commit() return { 'statusCode': 200, 'body': 'Data inserted successfully' } except Exception as e: print(f"Error: {str(e)}") return { 'statusCode': 500, 'body': str(e) } finally: if conn: cursor.close() conn.close()
配置Lambda层(解决pymysql依赖)
Lambda默认无pymysql库,需打包成层:
- 本地创建
python文件夹,执行pip install pymysql -t python/ - 将
python文件夹压缩为zip包 - 在Lambda控制台创建层,上传该zip包并关联到你的Lambda函数
配置AWS IoT Core规则
- 进入AWS IoT Core控制台创建规则
- 规则查询语句:
SELECT * FROM 'your/mqtt/topic'(替换为树莓派发布的MQTT主题) - 动作选择“Lambda”,关联你创建的Lambda函数
- 保存规则,IoT Core收到消息后将自动触发Lambda写入MySQL
方案二:EC2 订阅MQTT主题直接写入MySQL(适合初学者快速上手)
若觉得Lambda的VPC和层配置繁琐,可使用EC2实例运行Python脚本,直接订阅AWS IoT Core的MQTT主题,收到消息后写入MySQL。
前置准备
- EC2实例需具备互联网访问权限(默认配置即可)
- EC2安全组开放3306端口,允许访问MySQL数据库
- 在EC2上安装依赖:
pip install paho-mqtt pymysql - 从AWS IoT Core控制台下载设备证书(创建Thing后,下载证书、私钥、根CA证书)
Python 示例代码
import paho.mqtt.client as mqtt import pymysql import json # MySQL配置,替换为你的实际信息 MYSQL_CONFIG = { 'host': 'your-mysql-host', 'user': 'your-mysql-user', 'password': 'your-mysql-password', 'database': 'your-db-name', 'port': 3306 } # AWS IoT Core配置,替换为你的端点和证书路径 AWS_IOT_ENDPOINT = 'your-iot-endpoint-ats.iot.region.amazonaws.com' CA_CERT_PATH = './AmazonRootCA1.pem' DEVICE_CERT_PATH = './device-certificate.pem.crt' PRIVATE_KEY_PATH = './private.pem.key' MQTT_TOPIC = 'your/mqtt/topic' def on_connect(client, userdata, flags, rc): print(f"Connected with result code {rc}") client.subscribe(MQTT_TOPIC) def on_message(client, userdata, msg): payload = json.loads(msg.payload.decode()) print(f"Received message: {payload}") try: conn = pymysql.connect(**MYSQL_CONFIG) cursor = conn.cursor() insert_sql = """ INSERT INTO iot_data (timestamp, device_id, temperature) VALUES (%s, %s, %s) """ data = (payload['timestamp'], payload['device_id'], payload['temperature']) cursor.execute(insert_sql, data) conn.commit() print("Data inserted successfully") except Exception as e: print(f"Error inserting data: {str(e)}") finally: if conn: cursor.close() conn.close() # 初始化MQTT客户端 client = mqtt.Client() client.tls_set(CA_CERT_PATH, certfile=DEVICE_CERT_PATH, keyfile=PRIVATE_KEY_PATH) client.on_connect = on_connect client.on_message = on_message client.connect(AWS_IOT_ENDPOINT, 8883, 60) # 保持连接并处理消息 client.loop_forever()
运行脚本
将脚本和证书文件上传至EC2,执行python mqtt_subscribe_write_mysql.py,可使用nohup让脚本后台运行:nohup python mqtt_subscribe_write_mysql.py &
注意事项
- 需提前创建MySQL表,示例SQL:
CREATE TABLE iot_data ( id INT AUTO_INCREMENT PRIMARY KEY, timestamp DATETIME NOT NULL, device_id VARCHAR(50) NOT NULL, temperature FLOAT, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ); - 避免硬编码敏感信息,Lambda方案用Secrets Manager,EC2方案可使用环境变量
- 完善数据库连接的异常处理,避免程序崩溃
- 高并发场景下,Lambda可配置并发数,EC2方案可考虑引入消息队列缓冲请求
内容的提问来源于stack exchange,提问作者Wikinger749 A.D
相关产品推荐
相关产品推荐

