Phần 6: Mở rộng MQTT cho hệ thống giám sát chất lượng nước
Nâng cấp project thành hệ thống IoT bằng cách kết nối Wi-Fi, publish dữ liệu JSON lên MQTT và gửi alert khi phát hiện bất thường.
Phần 6: Mở rộng MQTT cho hệ thống giám sát chất lượng nước
Sau khi đọc cảm biến ổn định, bước tiếp theo là gửi dữ liệu realtime lên MQTT Broker. MQTT giúp ESP32 truyền dữ liệu nhẹ, dễ tích hợp với dashboard, backend, Node-RED, ThingsBoard hoặc Home Assistant.
6.1. Kiến trúc MQTT
[ESP32 Water Monitor]
→ MQTT Broker
→ Node-RED / Backend / Dashboard
→ Database / Alert
6.2. Topic đề xuất
iotlabs/water/water-monitor-001/telemetry
iotlabs/water/water-monitor-001/status
iotlabs/water/water-monitor-001/alert
iotlabs/water/water-monitor-001/heartbeat
6.3. Payload telemetry
{
"device_id": "water-monitor-001",
"temperature_c": 28.37,
"ph": 7.02,
"tds_ppm": 184.5,
"turbidity_percent": 2.2,
"status": "NORMAL",
"uptime_ms": 123456,
"wifi_rssi": -55
}
6.4. Cài thư viện MQTT
Trong Arduino IDE, cài thêm:
PubSubClientArduinoJson
6.5. Cấu hình Wi-Fi và MQTT
const char* WIFI_SSID = "YOUR_WIFI_NAME";
const char* WIFI_PASSWORD = "YOUR_WIFI_PASSWORD";
const char* MQTT_HOST = "192.168.1.50";
const int MQTT_PORT = 1883;
const char* MQTT_USERNAME = "";
const char* MQTT_PASSWORD = "";
const char* DEVICE_ID = "water-monitor-001";
6.6. Code hoàn chỉnh ESP32 + MQTT
#include <WiFi.h>
#include <PubSubClient.h>
#include <ArduinoJson.h>
#include <OneWire.h>
#include <DallasTemperature.h>
const char* WIFI_SSID = "YOUR_WIFI_NAME";
const char* WIFI_PASSWORD = "YOUR_WIFI_PASSWORD";
const char* MQTT_HOST = "192.168.1.50";
const int MQTT_PORT = 1883;
const char* MQTT_USERNAME = "";
const char* MQTT_PASSWORD = "";
const char* DEVICE_ID = "water-monitor-001";
String TOPIC_TELEMETRY;
String TOPIC_STATUS;
String TOPIC_ALERT;
String TOPIC_HEARTBEAT;
#define PIN_ONE_WIRE 4
#define PIN_TDS 34
#define PIN_PH 35
#define PIN_TURBIDITY 32
#define PIN_BUZZER 25
#define PIN_LED_GREEN 26
#define PIN_LED_RED 27
const float ADC_MAX = 4095.0;
const float ESP32_ADC_REF_VOLTAGE = 3.3;
const float TURBIDITY_DIVIDER_FACTOR = 1.5;
float PH7_VOLTAGE = 2.50;
float PH4_VOLTAGE = 3.00;
float CLEAR_WATER_VOLTAGE = 4.10;
float DIRTY_WATER_VOLTAGE = 2.50;
const float PH_MIN_NORMAL = 6.5;
const float PH_MAX_NORMAL = 8.5;
const float TDS_WARNING = 500.0;
const float TDS_DANGER = 1000.0;
const float TURBIDITY_WARNING = 50.0;
const float TURBIDITY_DANGER = 80.0;
const unsigned long TELEMETRY_INTERVAL_MS = 5000;
const unsigned long HEARTBEAT_INTERVAL_MS = 30000;
unsigned long lastTelemetryMs = 0;
unsigned long lastHeartbeatMs = 0;
String lastStatus = "";
WiFiClient espClient;
PubSubClient mqttClient(espClient);
OneWire oneWire(PIN_ONE_WIRE);
DallasTemperature tempSensor(&oneWire);
float readAdcVoltage(int pin) {
const int sampleCount = 30;
long total = 0;
for (int i = 0; i < sampleCount; i++) {
total += analogRead(pin);
delay(5);
}
float rawAverage = total / (float)sampleCount;
return rawAverage * ESP32_ADC_REF_VOLTAGE / ADC_MAX;
}
float readTemperatureC() {
tempSensor.requestTemperatures();
float temperatureC = tempSensor.getTempCByIndex(0);
if (temperatureC < -50 || temperatureC > 125) return NAN;
return temperatureC;
}
float calculateTds(float voltage, float temperatureC) {
if (isnan(temperatureC)) temperatureC = 25.0;
float compensationCoefficient = 1.0 + 0.02 * (temperatureC - 25.0);
float compensationVoltage = voltage / compensationCoefficient;
float tdsValue = (133.42 * compensationVoltage * compensationVoltage * compensationVoltage
- 255.86 * compensationVoltage * compensationVoltage
+ 857.39 * compensationVoltage) * 0.5;
if (tdsValue < 0) tdsValue = 0;
return tdsValue;
}
float calculatePh(float voltage) {
float slope = (7.0 - 4.0) / (PH7_VOLTAGE - PH4_VOLTAGE);
float intercept = 7.0 - slope * PH7_VOLTAGE;
return slope * voltage + intercept;
}
float calculateTurbidityPercent(float sensorVoltage) {
float percent = (CLEAR_WATER_VOLTAGE - sensorVoltage) * 100.0 / (CLEAR_WATER_VOLTAGE - DIRTY_WATER_VOLTAGE);
if (percent < 0) percent = 0;
if (percent > 100) percent = 100;
return percent;
}
String evaluateWaterStatus(float temperatureC, float ph, float tds, float turbidityPercent) {
if (isnan(temperatureC) || isnan(ph) || isnan(tds) || isnan(turbidityPercent)) return "SENSOR_ERROR";
if (ph < 4.0 || ph > 11.0) return "DANGER";
if (tds >= TDS_DANGER || turbidityPercent >= TURBIDITY_DANGER) return "DANGER";
if (ph < PH_MIN_NORMAL || ph > PH_MAX_NORMAL || tds >= TDS_WARNING || turbidityPercent >= TURBIDITY_WARNING) return "WARNING";
return "NORMAL";
}
void updateAlertOutput(String status) {
if (status == "NORMAL") {
digitalWrite(PIN_LED_GREEN, HIGH);
digitalWrite(PIN_LED_RED, LOW);
digitalWrite(PIN_BUZZER, LOW);
} else if (status == "WARNING") {
digitalWrite(PIN_LED_GREEN, LOW);
digitalWrite(PIN_LED_RED, HIGH);
digitalWrite(PIN_BUZZER, LOW);
} else {
digitalWrite(PIN_LED_GREEN, LOW);
digitalWrite(PIN_LED_RED, HIGH);
digitalWrite(PIN_BUZZER, HIGH);
}
}
void connectWiFi() {
if (WiFi.status() == WL_CONNECTED) return;
Serial.print("Connecting to Wi-Fi: ");
Serial.println(WIFI_SSID);
WiFi.mode(WIFI_STA);
WiFi.begin(WIFI_SSID, WIFI_PASSWORD);
int retry = 0;
while (WiFi.status() != WL_CONNECTED && retry < 30) {
delay(500);
Serial.print(".");
retry++;
}
Serial.println();
if (WiFi.status() == WL_CONNECTED) {
Serial.println("Wi-Fi connected");
Serial.print("IP address: "); Serial.println(WiFi.localIP());
Serial.print("RSSI: "); Serial.println(WiFi.RSSI());
} else {
Serial.println("Wi-Fi connection failed");
}
}
void setupMqttTopics() {
TOPIC_TELEMETRY = "iotlabs/water/" + String(DEVICE_ID) + "/telemetry";
TOPIC_STATUS = "iotlabs/water/" + String(DEVICE_ID) + "/status";
TOPIC_ALERT = "iotlabs/water/" + String(DEVICE_ID) + "/alert";
TOPIC_HEARTBEAT = "iotlabs/water/" + String(DEVICE_ID) + "/heartbeat";
}
void connectMqtt() {
if (mqttClient.connected()) return;
if (WiFi.status() != WL_CONNECTED) return;
Serial.print("Connecting to MQTT broker: ");
Serial.println(MQTT_HOST);
String clientId = "esp32-" + String(DEVICE_ID) + "-" + String(random(0xffff), HEX);
bool connected = false;
if (strlen(MQTT_USERNAME) > 0) {
connected = mqttClient.connect(clientId.c_str(), MQTT_USERNAME, MQTT_PASSWORD);
} else {
connected = mqttClient.connect(clientId.c_str());
}
if (connected) {
Serial.println("MQTT connected");
mqttClient.publish(TOPIC_STATUS.c_str(), "online", true);
} else {
Serial.print("MQTT connection failed, rc=");
Serial.println(mqttClient.state());
}
}
template <typename T>
T roundTo(T value, float factor) {
return round(value * factor) / factor;
}
bool publishJson(String topic, JsonDocument& doc, bool retained = false) {
char payload[512];
size_t length = serializeJson(doc, payload);
Serial.print("Publish topic: "); Serial.println(topic);
Serial.print("Payload: "); Serial.println(payload);
return mqttClient.publish(topic.c_str(), payload, length, retained);
}
void publishTelemetry(float temperatureC, float ph, float tds, float turbidityPercent, String status) {
StaticJsonDocument<512> doc;
doc["device_id"] = DEVICE_ID;
doc["temperature_c"] = round(temperatureC * 100) / 100.0;
doc["ph"] = round(ph * 100) / 100.0;
doc["tds_ppm"] = round(tds * 10) / 10.0;
doc["turbidity_percent"] = round(turbidityPercent * 10) / 10.0;
doc["status"] = status;
doc["uptime_ms"] = millis();
doc["wifi_rssi"] = WiFi.RSSI();
publishJson(TOPIC_TELEMETRY, doc, false);
mqttClient.publish(TOPIC_STATUS.c_str(), status.c_str(), true);
}
void publishAlert(float ph, float tds, float turbidityPercent, String status) {
StaticJsonDocument<512> doc;
doc["device_id"] = DEVICE_ID;
doc["level"] = status;
doc["message"] = "Water quality alert detected";
doc["ph"] = round(ph * 100) / 100.0;
doc["tds_ppm"] = round(tds * 10) / 10.0;
doc["turbidity_percent"] = round(turbidityPercent * 10) / 10.0;
doc["uptime_ms"] = millis();
publishJson(TOPIC_ALERT, doc, false);
}
void publishHeartbeat() {
StaticJsonDocument<256> doc;
doc["device_id"] = DEVICE_ID;
doc["status"] = "online";
doc["uptime_ms"] = millis();
doc["wifi_rssi"] = WiFi.RSSI();
doc["free_heap"] = ESP.getFreeHeap();
publishJson(TOPIC_HEARTBEAT, doc, false);
}
void setup() {
Serial.begin(115200);
delay(1000);
randomSeed(micros());
analogReadResolution(12);
analogSetAttenuation(ADC_11db);
tempSensor.begin();
pinMode(PIN_BUZZER, OUTPUT);
pinMode(PIN_LED_GREEN, OUTPUT);
pinMode(PIN_LED_RED, OUTPUT);
digitalWrite(PIN_BUZZER, LOW);
digitalWrite(PIN_LED_GREEN, LOW);
digitalWrite(PIN_LED_RED, LOW);
setupMqttTopics();
connectWiFi();
mqttClient.setServer(MQTT_HOST, MQTT_PORT);
mqttClient.setBufferSize(512);
connectMqtt();
Serial.println("System started");
Serial.println("ESP32 Water Quality Monitor + MQTT");
}
void loop() {
connectWiFi();
connectMqtt();
if (mqttClient.connected()) mqttClient.loop();
unsigned long now = millis();
if (now - lastTelemetryMs >= TELEMETRY_INTERVAL_MS) {
lastTelemetryMs = now;
float temperatureC = readTemperatureC();
float tdsVoltage = readAdcVoltage(PIN_TDS);
float phVoltage = readAdcVoltage(PIN_PH);
float turbidityEsp32Voltage = readAdcVoltage(PIN_TURBIDITY);
float turbiditySensorVoltage = turbidityEsp32Voltage * TURBIDITY_DIVIDER_FACTOR;
float tds = calculateTds(tdsVoltage, temperatureC);
float ph = calculatePh(phVoltage);
float turbidityPercent = calculateTurbidityPercent(turbiditySensorVoltage);
String status = evaluateWaterStatus(temperatureC, ph, tds, turbidityPercent);
updateAlertOutput(status);
Serial.println("====================================");
Serial.println("ESP32 Water Quality Monitor + MQTT");
Serial.print("Temperature: "); Serial.print(temperatureC, 2); Serial.println(" °C");
Serial.print("TDS: "); Serial.print(tds, 1); Serial.println(" ppm");
Serial.print("pH: "); Serial.println(ph, 2);
Serial.print("Turbidity Relative: "); Serial.print(turbidityPercent, 1); Serial.println(" %");
Serial.print("Status: "); Serial.println(status);
Serial.print("Wi-Fi: "); Serial.println(WiFi.status() == WL_CONNECTED ? "CONNECTED" : "DISCONNECTED");
Serial.print("MQTT: "); Serial.println(mqttClient.connected() ? "CONNECTED" : "DISCONNECTED");
Serial.println("====================================");
if (mqttClient.connected()) {
publishTelemetry(temperatureC, ph, tds, turbidityPercent, status);
if (status != "NORMAL" && status != lastStatus) {
publishAlert(ph, tds, turbidityPercent, status);
}
}
lastStatus = status;
}
if (now - lastHeartbeatMs >= HEARTBEAT_INTERVAL_MS) {
lastHeartbeatMs = now;
if (mqttClient.connected()) publishHeartbeat();
}
}
6.7. Kết quả mong đợi
Wi-Fi connected
MQTT connected
Publish topic: iotlabs/water/water-monitor-001/telemetry
Payload: {"device_id":"water-monitor-001","temperature_c":28.37,"ph":7.02,"tds_ppm":184.5,"turbidity_percent":2.2,"status":"NORMAL"}