Este tutorial veremos cómo realizar la integración de MQTT e InfluxDB usando N8N. De este modo, podrás almacenar los datos de dispositivos IoT que utilizan MQTT en un bucket de InfluxDB Cloud.

Figura 1 – Sistema propuesto

Descripción General del Workflow

El workflow analizado está diseñado para recibir mensajes MQTT de dispositivos IoT, transformar los datos al formato requerido por InfluxDB (protocolo de línea) y enviarlos a InfluxDB Cloud para su almacenamiento.

El workflow consta de tres nodos principales:

  • MQTT Trigger: Recibe mensajes de un broker MQTT.
  • Code: Transforma los datos JSON en formato de protocolo de línea InfluxDB.
  • HTTP Request: Envía los datos transformados a InfluxDB Cloud.

Figura 2 – Workflow

Prerrequisitos

Antes de implementar este workflow, necesitarás tener configurados los siguientes elementos:

  • Una instancia de N8N funcionando (local o en la nube).
  • Acceso a un broker MQTT en la nube (como EMQX, Mosquitto, HiveMQ, etc.).
  • Una cuenta en InfluxDB Cloud con un bucket configurado.
  • Dispositivos IoT o simuladores que publiquen datos en formato JSON.

Implementación Paso a Paso

A continuación se detallan los pasos para implementar el workflow completo.

1: Configuración del Broker MQTT

  1. Asegúrate de tener un broker MQTT configurado y funcionando correctamente.
  2. Anota la dirección del broker, puerto, nombre de usuario y contraseña, ya que los necesitarás más adelante.

2: Configuración de InfluxDB Cloud

  1. Crea una cuenta en InfluxDB Cloud si aún no tienes una.
  2. Configura un bucket para almacenar tus datos IoT.
  3. Genera un token de API con permisos de escritura para el bucket.
  4. Anota la URL de tu organización, ID de organización, nombre del bucket y token API.

3: Creación del Workflow en N8N

Configuración del Nodo MQTT Trigger

Sigue estos pasos para configurar el nodo MQTT (trigger) en N8N:

  1. Crea un nuevo workflow en N8N.
  2. Agrega un nodo «MQTT Trigger» al canvas.
    • Configura las credenciales del broker MQTT:
      • Protocolo (mqtt/mqtts): En este caso utilizaremos mqtts, que provee encriptación de los datos. Sin embargo, puedes utilizar la opción que se ajuste a tus necesidades.
      • Host: Aquí debes ingresar el nombre del host de la instancia del broker MQTT.
      • Puerto: El puerto por defecto para conexiones MQTT sin encriptación es 1883, mientras que para conexiones encriptadas es 8883. Verifica esto en la configuración del broker que estés utilizando.
      • Usuario y contraseña: Si tu configuración del cliente MQTT utiliza usuario y contraseña, completa estos valores en el nodo N8N.
      • Client ID: Es recomendable dejar el Client ID vacío para que N8N genere uno automáticamente cada vez que se conecte al broker.
      • Topic: Configura el topic MQTT para suscribirse (ejemplo: «buenosaires/hotel01/room103»).

Configuración del Nodo Code para Transformación

  1. Agrega un nodo «Code» después del MQTT Trigger.
  2. Selecciona el modo «JavaScript». El código JavaScript que vas a necesitar dependerá del formato que tenga el payload del mensaje MQTT.
    Por ejemplo, si el payload es como el siguiente, se puede utilizar la función JavaScript propuesta más abajo para construir la línea de protocolo que necesaria para escribir los datos en el bucket de InfluxDB.

Payload del sensor.

{
«timestamp»: «2025-05-30T16:30:00Z»,
«device_id»: «sensor_001»,
«temperature»: 25.3,
«humidity»: 60.8,
«pressure»: 1013.25,
«location»: «zona_a»,
«battery_level»: 87
}

Función JavaScript para dar formato de «line protocol», que es utilizado por InfluxDB.

for (const item of $input.all()) {
  const payload = JSON.parse(item.json.message);

  // Conversión de timestamp
  const timestampMs = Date.parse(payload.timestamp);
  const timestampNs = timestampMs * 1e6;

  // --- TAGS ---
  const tags = [];
  // Valida y sanitiza cada tag
  if (payload.device_id?.trim()) tags.push(`device_id=${payload.device_id.trim()}`);
  if (payload.location?.trim()) tags.push(`location=${payload.location.trim()}`);
//  if (item.json.topic?.trim()) tags.push(`topic=${item.json.topic.trim()}`);

  // --- FIELDS ---
  const fields = [];
  // Valida que los campos sean numéricos
  if (typeof payload.temperature === 'number') fields.push(`temperature=${payload.temperature}`);
  if (typeof payload.humidity === 'number') fields.push(`humidity=${payload.humidity}`);
  if (typeof payload.pressure === 'number') fields.push(`pressure=${payload.pressure}`);
  if (typeof payload.battery_level === 'number') fields.push(`battery_level=${payload.battery_level}`);

  // --- VALIDACIÓN CRÍTICA ---
  if (tags.length === 0) throw new Error('Error: Todos los tags están vacíos');
  if (fields.length === 0) throw new Error('Error: Todos los fields están vacíos');

  // --- CONSTRUCCIÓN DE LA LÍNEA ---
  const measurement = 'iot_sensors';
  const tagsSection = tags.join(','); // Ej: "device_id=sensor_001,location=zona_a"
  const fieldsSection = fields.join(','); // Ej: "temperature=25.3,humidity=60.8"

  // Combina todo evitando comas extras
  const line = `${measurement},${tagsSection} ${fieldsSection} ${timestampNs}`;

  return [{ json: { line } }];
}

El código realiza las siguientes operaciones:

  • Parsea el mensaje JSON recibido.
  • Convierte el timestamp ISO8601 a nanosegundos para InfluxDB.
  • Organiza los datos en tags (metadatos indexados) y fields (valores de medición).
  • Valida que existan tags y fields.
  • Construye la línea de protocolo InfluxDB con el formato adecuado.

Configuración del Nodo HTTP Request

Finalmente, resta configurar el nodo HTTP Request para enviar los datos al bucket InfluxDB.

  1. Agrega un nodo «HTTP Request» después del nodo Code.
  2. Configura el método como POST.
  3. Establece la URL del endpoint de escritura de InfluxDB Cloud. Ten en cuenta la región en la que está corriendo la instancia de InfluxDB y reemplaza el valor en la siguiente URL: https://[REGION].aws.cloud2.influxdata.com/api/v2/write
  4. Configura los parámetros de query (ver Figura 3):
    • org ID de tu organizaciónbucket.
    • Nombre de tu bucket.
    • Precisión: ns (nanosegundos).
  5. Configura los headers (ver Figura 4):
    • Authorization: Token [tu-token-de-API]
    • Content-Type: text/plain; charset=utf-8
    • Accept: application/json
  6. En el body, usa una expresión para referenciar la línea generada en el bloque de código (ver código JavaScript y Figura 4): ={{ $json.line }}

Figura 3 – Configuración de URL y parámetros de query.

Figura 4 – Configuración de headers y body.

Problemas comunes en la implentación

Durante la implementación de workflows pueden surgir problemas relacionados al formato de los datos, incorrecto uso de las APIs, error en la autenticación, etc.
Veamos algunos ejemplos que pueden ocurrir para este workflow en particular.

Problemas con el Formato del Protocolo de Línea

El error más común es recibir un código 400 con mensajes como «empty tag name» o «invalid line protocol». Esto suele ocurrir porque:

  • Hay una coma entre la sección de tags y fields (debe haber un espacio, no una coma).
  • Algún tag tiene nombre o valor vacío.
  • El formato del timestamp es incorrecto.

La solución implementada en el código de ejemplo verifica y valida todos estos aspectos antes de enviar los datos. Puedes usar este código como referencia para crear tu propio formateador.

Problemas de Conexión MQTT

Si el nodo MQTT Trigger no recibe mensajes, verifica:

  • La conectividad al broker.
  • La configuración correcta del topic.
  • Que no haya conflictos de conexiones múltiples (dejar el Client ID vacío ayuda).

Mejores Prácticas

Para entornos de producción robustos, considera estas recomendaciones:

  • Manejo de errores: Agrega nodos IF para validar datos y filtrar mensajes malformados
  • Escalabilidad: Implementa batching de escrituras agrupando múltiples mediciones en una sola solicitud para alto volumen de datos.
  • Seguridad: Utiliza el sistema de gestión de credenciales de n8n y habilita TLS/SSL en las conexiones MQTT cuando sea posible.
  • Monitoreo: Configura notificaciones por email o webhook cuando ocurran errores críticos.

Conclusión

Este workflow proporciona una solución escalable y robusta para integrar dispositivos IoT que usan MQTT con InfluxDB Cloud. La transformación optimizada de datos garantiza la confiabilidad y eficiencia en la ingesta de datos de series temporales, permitiendo un análisis en tiempo real de tus datos IoT.

La combinación de N8N con MQTT e InfluxDB ofrece una arquitectura flexible que puede adaptarse a diversos casos de uso de IoT, desde monitoreo ambiental hasta automatización industrial.


0 comentarios

Deja una respuesta

Marcador de posición del avatar

Tu dirección de correo electrónico no será publicada. Los campos obligatorios están marcados con *

Este sitio usa Akismet para reducir el spam. Aprende cómo se procesan los datos de tus comentarios.