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

如何用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库,需打包成层:

  1. 本地创建python文件夹,执行pip install pymysql -t python/
  2. 将python文件夹压缩为zip包
  3. 在Lambda控制台创建层,上传该zip包并关联到你的Lambda函数

配置AWS IoT Core规则

  1. 进入AWS IoT Core控制台创建规则
  2. 规则查询语句:SELECT * FROM 'your/mqtt/topic'(替换为树莓派发布的MQTT主题)
  3. 动作选择“Lambda”,关联你创建的Lambda函数
  4. 保存规则,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 05:10:37