不會爆炸的工廠——ATLANTIS壓力傳送器×AWS IoT防爆監控完整指南
不會爆炸的工廠——ATLANTIS × AWS IoT 工程師實戰指南
AWS IoT軟體工程師必讀|邊界AI推論 × MQTT 訊息路由 × 預測性維護 × 成本優化 || 昶特資訊部工程師 撰
🏗️ 系統架構五層模型
完整數據流架構:
↓
【Layer 2】AWS Greengrass 邊界閘道 + 本地 Lambda + TensorFlow Lite AI
↓
【Layer 3】AWS IoT Core MQTT Broker (本地 + 雲端)
↓
【Layer 4】IoT Rules + Lambda + DynamoDB + Timestream
↓
【Layer 5】CloudWatch Dashboard + SNS 告警 + QuickSight 可視化
🔌 Layer 1:感測器通訊轉換
1.1 HART 與 Modbus 協議支持
ATLANTIS SDPT-3100 智能型壓力傳送器支援兩種工業標準協議。Greengrass 邊界須負責將兩種協議統一轉換為 MQTT JSON 格式。
核心轉換邏輯:
- HART (1200 baud): Preamble + Delimiter + Address + Command 3 + Status + PV (4字節 IEEE 754浮點) + Checksum
- Modbus RTU (9600 baud): Slave ID + Function Code 3/4 + 寄存器地址 + 寄存器數量 + CRC16
- MQTT 輸出:
{"pressure_bar": 45.67, "timestamp": 1234567890, "status": "OK"}
💾 完整代碼:HART Serial → MQTT 轉換
Python:HART 讀取與 AWS IoT 發送
import json
import time
from AWSIoTPythonSDK.MQTTLib import AWSIoTMQTTClient
import serial
import struct
HART_PORT = '/dev/ttyUSB0'
HART_BAUDRATE = 1200
CLIENT_ID = "greengrass-sdpt3100-reader"
THING_NAME = "factory-pressure-monitor-01"
class HARTReader:
def __init__(self, port, baudrate):
self.serial = serial.Serial(port, baudrate, timeout=1.0)
self.frame_buffer = bytearray()
def read_hart_frame(self):
"""HART幀結構"""
self.frame_buffer.clear()
while True:
byte = self.serial.read(1)
if not byte:
raise TimeoutError("HART讀取超時")
if byte == b'\x02' or byte == b'\x82':
self.frame_buffer.append(byte[0])
break
while len(self.frame_buffer) < 20:
byte = self.serial.read(1)
if not byte:
break
self.frame_buffer.append(byte[0])
return self.parse_hart_response()
def parse_hart_response(self):
"""解析HART回應,提取壓力值"""
if len(self.frame_buffer) < 10:
return None
status = self.frame_buffer[2]
response_code = self.frame_buffer[5]
if response_code != 0x00:
return None
# 提取PV (4字節IEEE 754浮點)
pv_bytes = self.frame_buffer[6:10]
pv_value = struct.unpack('>f', bytes(pv_bytes))[0]
device_status = 'FAULT' if (status & 0x80) else 'OK'
return {
'pv': round(pv_value, 2),
'device_status': device_status
}
class IoTPublisher:
def __init__(self, client_id):
self.client = AWSIoTMQTTClient(client_id)
self.client.configureEndpoint(
"a1b2c3d4e5f6g7h8.iot.ap-northeast-1.amazonaws.com", 8883)
self.client.configureCredentials(
"/greengrass/certs/AmazonRootCA1.pem",
"/greengrass/certs/private.key",
"/greengrass/certs/certificate.pem.crt")
self.connected = False
def connect(self):
self.client.connect()
self.connected = True
def publish_pressure_data(self, pressure_data, sensor_id):
payload = {
"state": {
"reported": {
"sensor_id": sensor_id,
"timestamp": int(time.time()),
"pressure_bar": pressure_data['pv'],
"device_status": pressure_data['device_status']
}
}
}
topic = f"$aws/things/{THING_NAME}/shadow/update"
self.client.publish(topic, json.dumps(payload), 1)
if __name__ == "__main__":
hart_reader = HARTReader(HART_PORT, HART_BAUDRATE)
iot_publisher = IoTPublisher(CLIENT_ID)
iot_publisher.connect()
while True:
try:
pressure_data = hart_reader.read_hart_frame()
if pressure_data:
iot_publisher.publish_pressure_data(pressure_data, "SDPT3100_001")
time.sleep(5)
except TimeoutError:
print("HART通訊超時,重試中...")
time.sleep(2)
💾 完整代碼:Modbus RTU → MQTT 轉換
Python:RS-485 Modbus 讀取與發送
import minimalmodbus
import json
import time
from AWSIoTPythonSDK.MQTTLib import AWSIoTMQTTClient
MODBUS_PORT = '/dev/ttyUSB1'
MODBUS_SLAVE_ID = 1
MODBUS_BAUDRATE = 9600
class ModbusToMQTT:
def __init__(self):
self.instrument = minimalmodbus.Instrument(MODBUS_PORT, MODBUS_SLAVE_ID)
self.instrument.serial.baudrate = MODBUS_BAUDRATE
self.instrument.serial.timeout = 1.0
self.mqtt_client = AWSIoTMQTTClient("greengrass-dptx-reader")
self.mqtt_client.configureEndpoint(
"a1b2c3d4e5f6g7h8.iot.ap-northeast-1.amazonaws.com", 8883)
self.mqtt_client.configureCredentials(
"/greengrass/certs/AmazonRootCA1.pem",
"/greengrass/certs/private.key",
"/greengrass/certs/certificate.pem.crt")
def read_pressure_modbus(self):
"""讀取DPTX壓力值 (Modbus RTU)"""
try:
registers = self.instrument.read_registers(
registeraddress=0,
number_of_registers=3,
functioncode=3)
pressure_value = registers[0] + registers[1]/10000.0
device_status = 'OK' if (registers[2] & 0x0001) else 'ERROR'
return {
'pressure_bar': round(pressure_value, 2),
'device_status': device_status,
'timestamp': int(time.time())
}
except Exception as e:
print(f"Modbus讀取失敗: {str(e)}")
return None
def publish_to_mqtt(self, sensor_id, pressure_data):
payload = {
"state": {
"reported": {
"sensor_id": sensor_id,
"pressure_bar": pressure_data['pressure_bar'],
"device_status": pressure_data['device_status'],
"timestamp": pressure_data['timestamp'],
"protocol": "Modbus RTU"
}
}
}
topic = "$aws/things/factory-pressure-monitor-02/shadow/update"
self.mqtt_client.publish(topic, json.dumps(payload), 1)
if __name__ == "__main__":
converter = ModbusToMQTT()
converter.mqtt_client.connect()
while True:
data = converter.read_pressure_modbus()
if data:
converter.publish_to_mqtt("DPTX_001", data)
time.sleep(5)
📡 Layer 2:Greengrass 邊界AI推論
2.1 架構概覽
AWS Greengrass 在工廠邊界設備運行,負責本地 MQTT Broker、Lambda 函數、TensorFlow Lite 異常檢測模型、本地 DynamoDB 和離線消息隊列。
🏗️ Greengrass 部署架構圖
│ 工廠網路 (本地 LAN) │
├─────────────────────────────────────────┤
│ ┌───────────────────────────────────┐ │
│ │ Greengrass Core (Ubuntu 20.04) │ │
│ │ │ │
│ │ ├─ Local MQTT Broker (埠8883) │ │
│ │ ├─ Lambda #1: HART解析 (5秒週期) │ │
│ │ ├─ Lambda #2: 異常檢測AI (<100ms) │ │
│ │ ├─ Lambda #3: 雲端同步 │ │
│ │ └─ Local DynamoDB (1小時歷史) │ │
│ │ │ │
│ └───────────────────────────────────┘ │
│ ↑ ↓ HART/RS-485 (50個感測器) │
└─────────────────────────────────────────┘
↓ MQTT over TLS (網際網路)
┌─────────────────────────────────────────┐
│ AWS 雲端 │
│ IoT Core → Rules → Lambda → DynamoDB │
└─────────────────────────────────────────┘
2.2 異常檢測 AI 模型
TensorFlow Lite 異常檢測模型特點:
- 模型大小:1.2 MB(量化版本)
- 推論延遲:15 ms(相比雲端 150-500ms 快 10 倍)
- 準度:94-96% 異常檢測準確率
- 特徵提取:15 個特徵(統計、頻域、時間序列)
- 異常閾值:概率 > 0.7 觸發警報
💾 TensorFlow Lite 異常檢測 Lambda 函數
Python:Greengrass Lambda 異常檢測
import json
import numpy as np
import tensorflow as tf
import time
import logging
logger = logging.getLogger()
MODEL_PATH = "/greengrass/ml/pressure_anomaly_model.tflite"
interpreter = tf.lite.Interpreter(model_path=MODEL_PATH)
interpreter.allocate_tensors()
input_details = interpreter.get_input_details()
output_details = interpreter.get_output_details()
class PressureAnomalyDetector:
def __init__(self):
self.pressure_history = []
self.max_history_size = 60
self.anomaly_threshold = 0.7
self.normal_range = (0, 100)
self.rate_of_change_limit = 5.0
def extract_features(self, pressure_values):
"""從壓力序列提取15個特徵"""
if len(pressure_values) < 5:
return None
arr = np.array(pressure_values, dtype=np.float32)
mean_val = float(np.mean(arr))
std_val = float(np.std(arr))
min_val = float(np.min(arr))
max_val = float(np.max(arr))
median_val = float(np.median(arr))
diff1 = np.diff(arr)
diff_mean = float(np.mean(np.abs(diff1)))
diff_std = float(np.std(diff1))
diff2 = np.diff(diff1)
accel_mean = float(np.mean(np.abs(diff2))) if len(diff2) > 0 else 0.0
fft_vals = np.abs(np.fft.fft(arr))
fft_peak = float(np.max(fft_vals[1:len(fft_vals)//2]))
z_scores = np.abs((arr - mean_val) / (std_val + 1e-6))
outlier_count = float(np.sum(z_scores > 3.0))
extreme_ratio = float(np.sum((arr < self.normal_range[0]) |
(arr > self.normal_range[1])) / len(arr))
features = np.array([
mean_val, std_val, min_val, max_val, median_val,
diff_mean, diff_std, accel_mean,
fft_peak,
outlier_count, extreme_ratio,
mean_val * std_val,
diff_mean / (std_val + 1e-6),
max_val - min_val,
float(len(arr))
], dtype=np.float32)
return features
def predict_anomaly(self, pressure_value):
"""使用TF Lite模型預測異常概率"""
self.pressure_history.append(pressure_value)
if len(self.pressure_history) > self.max_history_size:
self.pressure_history.pop(0)
if len(self.pressure_history) < 10:
return {
'anomaly_probability': 0.0,
'status': 'INSUFFICIENT_DATA'
}
# 急劇變化檢查
if len(self.pressure_history) >= 2:
rate_of_change = abs(self.pressure_history[-1] -
self.pressure_history[-2])
if rate_of_change > self.rate_of_change_limit:
return {
'anomaly_probability': 0.95,
'status': 'RAPID_CHANGE_DETECTED'
}
# TensorFlow Lite推論
try:
features = self.extract_features(self.pressure_history[-10:])
if features is None:
return {'anomaly_probability': 0.0, 'status': 'ERROR'}
features_reshaped = features.reshape((1, -1))
interpreter.set_tensor(input_details[0]['index'], features_reshaped)
interpreter.invoke()
output_data = interpreter.get_tensor(output_details[0]['index'])
anomaly_prob = float(output_data[0][0])
if anomaly_prob > self.anomaly_threshold:
status = 'ANOMALY_DETECTED'
elif anomaly_prob > 0.3:
status = 'WARNING'
else:
status = 'NORMAL'
return {
'anomaly_probability': round(anomaly_prob, 4),
'status': status
}
except Exception as e:
logger.error(f"推論錯誤: {str(e)}")
return {'anomaly_probability': 0.0, 'status': 'INFERENCE_ERROR'}
def lambda_handler(event, context):
"""Greengrass Lambda入口點"""
detector = PressureAnomalyDetector()
sensor_id = event.get('sensor_id')
pressure = event.get('pressure_bar')
timestamp = event.get('timestamp', int(time.time()))
if not pressure or not sensor_id:
return {'statusCode': 400}
anomaly_result = detector.predict_anomaly(pressure)
alert_action = 'NO_ACTION'
alert_level = 'NONE'
if anomaly_result['status'] == 'ANOMALY_DETECTED':
alert_action = 'TRIGGER_ALERT'
alert_level = 'HIGH'
elif anomaly_result['status'] == 'WARNING':
alert_action = 'MONITOR'
alert_level = 'MEDIUM'
response = {
'sensor_id': sensor_id,
'pressure_bar': pressure,
'anomaly_probability': anomaly_result['anomaly_probability'],
'prediction_status': anomaly_result['status'],
'alert_action': alert_action,
'alert_level': alert_level
}
logger.info(f"異常檢測結果: {json.dumps(response)}")
return response
☁️ Layer 3 & 4:AWS IoT Core + 雲端處理
3.1 MQTT 主題架構設計
階層化主題結構:
factory/plant/{plant_id}/sensor/{sensor_type}/{sensor_id}/data
factory/plant/shimane-01/sensor/pressure/SDPT3100_001/data— 壓力數據factory/plant/shimane-01/alerts/anomaly/SDPT3100_001/critical— 異常警報factory/plant/shimane-01/diagnostics/device-health/SDPT3100_001— 診斷訊息
3.2 AWS IoT Rules(訊息路由)
💾 IoT Rule #1:壓力數據 → DynamoDB
SELECT
sensor_id,
pressure_bar,
device_status,
timestamp,
timestamp() as server_timestamp
FROM 'factory/plant/shimane-01/sensor/pressure/+/data'
WHERE pressure_bar > -1
ACTION:
DynamoDB:
RoleArn: arn:aws:iam::123456789012:role/iot-dynamodb-role
TableName: FactoryPressureData
HashKeyValue: ${sensor_id}
RangeKeyValue: ${timestamp}
Item:
pressure_bar
device_status
temperature_c
💾 IoT Rule #2:異常警報 → SNS 通知
SELECT
sensor_id,
anomaly_probability,
alert_level,
pressure_bar,
timestamp
FROM 'factory/plant/shimane-01/alerts/anomaly/+/critical'
WHERE anomaly_probability > 0.85
ACTION:
Sns:
RoleArn: arn:aws:iam::123456789012:role/iot-sns-role
TargetArn: arn:aws:sns:ap-northeast-1:123456789012:critical-alerts
MessageFormat: JSON
💾 雲端 Lambda:設備健康評分 (0~100)
import json
import boto3
import numpy as np
from datetime import datetime, timedelta
dynamodb = boto3.resource('dynamodb')
cloudwatch = boto3.client('cloudwatch')
def lambda_handler(event, context):
"""
分析過去1小時壓力數據,計算設備健康評分
評分維度:
1. 數據完整性 (0~25分)
2. 壓力穩定性 (0~25分)
3. 異常事件 (0~30分)
4. 設備狀態 (0~20分)
"""
table = dynamodb.Table('FactoryPressureData')
one_hour_ago = int((datetime.now() - timedelta(hours=1)).timestamp())
response = table.query(
KeyConditionExpression='sensor_id = :sid AND #ts > :ts',
ExpressionAttributeNames={'#ts': 'timestamp'},
ExpressionAttributeValues={
':sid': 'SDPT3100_001',
':ts': one_hour_ago
}
)
items = response.get('Items', [])
if not items:
return {'statusCode': 404}
pressures = [float(item['pressure_bar']) for item in items]
# 計算各維度評分
data_completeness_score = 25 * min(len(items) / 60, 1.0)
mean_p = np.mean(pressures)
std_p = np.std(pressures)
cv = std_p / (mean_p + 1e-6) if mean_p > 0 else 0
stability_score = 25 * max(0, 1 - cv / 0.15)
anomaly_count = sum(1 for item in items
if item.get('anomaly_probability', 0) > 0.7)
anomaly_score = 30 * max(0, 1 - anomaly_count / 10)
fault_count = sum(1 for item in items
if item.get('device_status') == 'FAULT')
device_score = 20 * max(0, 1 - fault_count / 5)
total_health_score = (data_completeness_score + stability_score +
anomaly_score + device_score)
cloudwatch.put_metric_data(
Namespace='AtlantisFactory',
MetricData=[{
'MetricName': 'DeviceHealthScore',
'Value': total_health_score,
'Dimensions': [
{'Name': 'SensorID', 'Value': 'SDPT3100_001'},
{'Name': 'Plant', 'Value': 'shimane-01'}
]
}]
)
return {
'statusCode': 200,
'health_score': round(total_health_score, 2),
'data_points_analyzed': len(items)
}
💰 成本優化與 ROI 分析
4.1 AWS IoT 成本構成(50 個感測器/月)
| 成本項目 | 計費方式 | 純雲端 | Greengrass優化 | 節省 |
|---|---|---|---|---|
| IoT Core 連接 | $0.10/百萬連線-秒 | $129.6 | $65 | 50%↓ |
| MQTT 訊息 | $1/百萬訊息 | $216 | $65 | 70%↓ |
| DynamoDB 寫入 | $1.25/百萬單位 | $150 | $45 | 70%↓ |
| Lambda 執行 | $0.20/百萬調用 | $45 | $5 | 89%↓ |
| 合計 | $631 | $215 | 66%↓ |
❓ 工程師常見問題 (FAQ)
1️⃣ Greengrass 無網路環境如何運作?離線緩衝機制?
Greengrass v2 內置本地 MQTT Broker 和離線消息隊列。當網際網路斷開時:
- 本地 Lambda 繼續執行異常檢測(讀取本地 DynamoDB 歷史)
- MQTT 訊息自動緩衝到本地磁碟(可存儲 24h+ 數據)
- 網路恢復時自動同步
- 配置在
/greengrass/config/config.yaml的mqttBroker.persistence
2️⃣ TensorFlow Lite 模型大小與推論速度的權衡?
三種配置方案:
- 完整模型 (5MB): 準度 97%,推論 50ms,啟動 2 秒
- 量化模型 (1.2MB): 準度 94-96%,推論 15ms,啟動 0.5 秒 ✅ 推薦
- 極端優化 (200KB): 準度 89%,推論 5ms(精度不足,不建議)
3️⃣ MQTT 主題設計最佳實踐?避免衝突?
階層化命名規則:
factory/plant/{plant_id}/sensor/{sensor_type}/{sensor_id}/data— 數據發佈factory/plant/{plant_id}/control/{device_id}/commands— 控制指令$aws/things/{thing_name}/shadow/update— Device Shadow
避免衝突: 使用 UUID 或 MAC 地址作為 sensor_id,禁止發佈主題用通配符。
4️⃣ Greengrass 記憶體管理:監控防止 OOM?
監控 Greengrass 進程記憶體使用,配置 Java 堆大小限制。防止 OOM:啟用 Lambda 記憶體限制 (256MB)、設置 DynamoDB TTL、定期重啟 Greengrass 服務。
5️⃣ 從 Greengrass Lambda 讀取本地 DynamoDB?
使用本地 DynamoDB 端點 http://localhost:8000,執行查詢操作並提取壓力值用於 AI 異常檢測模型輸入。
6️⃣ HART 通訊超時和重試機制?
HART 是低速通訊(1200 baud),容易超時。實現健壯重試邏輯:若讀取失敗,退避 0.5 秒後重試,最多 3 次。若全部失敗返回上次有效值或 NaN。
7️⃣ AWS IoT Core 訊息吞吐量限制?
默認限制:每個連線 1,000 訊息/秒。50 感測器方案每秒 50 訊息(安全)。優化策略:邊界聚合 5 秒後發送(訊息量 ↓ 5 倍)、使用 MQTT QoS 1、啟用訊息壓縮。
8️⃣ 傳感器時鐘不同步問題?
在 Greengrass 端使用 timestamp() 覆蓋感測器 timestamp,定期同步傳感器時鐘(via HART 指令),使用 AWS Timestream 自動時鐘校正功能。
9️⃣ 灰度部署新模型?
部署新異常檢測模型時進行灰度發佈:10% 設備運行 1 週 → 50% 設備運行 1 週 → 100% 部署。若新模型性能下降 5% 以上,自動回滾到舊版本。
🔟 與現有 SCADA 系統整合(OPC UA)?
Greengrass 可作為 OPC UA 橋接器,讀取現有 SCADA 服務器數據並轉發到 AWS IoT Core,實現新舊系統的無縫融合。
✅ 實施檢查清單
- ☐ Greengrass Core 2.0 安裝完畢(ARM/x86 支援)
- ☐ HART/RS-485 驅動及 Serial 埠訪問權限
- ☐ AWS IoT Thing 和 Certificate 建立完成
- ☐ TensorFlow Lite 異常檢測模型上傳到 Greengrass
- ☐ DynamoDB 表及 TTL 配置(自動刪除 30 天舊數據)
- ☐ AWS IoT Rules 三條規則部署(DynamoDB/Timestream/SNS)
- ☐ CloudWatch 儀錶板 + 告警設置
- ☐ 負載測試:50 感測器同時連線無丟包
- ☐ 離線測試:網路斷開後邊界 AI 仍能推論
- ☐ 安全審計:IAM 權限最小化、TLS v1.2+、證書有效期