Part 6: Extending with MQTT for the Water Quality Monitoring System
Upgrades the project into a full IoT system by connecting to WiFi, publishing JSON data over MQTT, and sending alerts when an anomaly is detected.
Part 6: Extending with MQTT for the Water Quality Monitoring System
Once sensor readings are stable, the next step is sending real-time data to an MQTT broker. MQTT lets the ESP32 transmit data lightweight and easily integrate with a dashboard, backend, Node-RED, ThingsBoard, or Home Assistant.
6.1. MQTT Architecture
[ESP32 Water Monitor]
→ MQTT Broker
→ Node-RED / Backend / Dashboard
→ Database / Alert
6.2. Suggested Topics
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. Telemetry Payload
{
"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. Installing MQTT Libraries
In the Arduino IDE, install additionally:
PubSubClientArduinoJson
6.5. WiFi and MQTT Configuration
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. Full ESP32 + MQTT Code
#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. Expected Result
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"}