移至主內容

不會爆炸的工廠——ATLANTIS壓力傳送器×AWS IoT防爆監控完整指南

不會爆炸的工廠——ATLANTIS × AWS IoT 工程師實戰指南

AWS IoT軟體工程師必讀|邊界AI推論 × MQTT 訊息路由 × 預測性維護 × 成本優化 || 昶特資訊部工程師 撰

💡 提示: 本頁面包含多個可展開區域(帶有 ▶ 箭頭)。點擊任何帶有箭頭的欄位即可展開完整的代碼、架構圖或詳細內容。點擊再次關閉。

🏗️ 系統架構五層模型

完整數據流架構:

【Layer 1】ATLANTIS壓力傳送器 (HART/RS-485)
      ↓
【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$6550%↓
MQTT 訊息$1/百萬訊息$216$6570%↓
DynamoDB 寫入$1.25/百萬單位$150$4570%↓
Lambda 執行$0.20/百萬調用$45$589%↓
合計 $631$21566%↓
550萬
初期投資(50個感測器 + Greengrass 主機 + AWS)
1,850萬
首年節省(避免停機 + 維保成本降低)
3.6個月
投資回收期

❓ 工程師常見問題 (FAQ)

1️⃣ Greengrass 無網路環境如何運作?離線緩衝機制?

Greengrass v2 內置本地 MQTT Broker 和離線消息隊列。當網際網路斷開時:

  • 本地 Lambda 繼續執行異常檢測(讀取本地 DynamoDB 歷史)
  • MQTT 訊息自動緩衝到本地磁碟(可存儲 24h+ 數據)
  • 網路恢復時自動同步
  • 配置在 /greengrass/config/config.yamlmqttBroker.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+、證書有效期