Estoy escribiendo un script node.js que usa mqtt para recibir mensajes publicados en un solo tema. Quiero ejecutar una función que tome un dato del mensaje y cree una lista de temas a los que se suscribe el cliente. Tengo el código de trabajo como se indica a continuación.
El problema es que tendré miles de dispositivos que enviarán mensajes al tema: "tema/información" y algunos cientos podrían ser mensajes simultáneos. Aunque los mensajes se pueden recibir con la configuración de QoS en MQTT, la ejecución en serie no es muy eficiente ya que tengo un bucle for en ejecución, así como una inserción de base de datos (no en el código, solo una línea comentada) dentro de la función. Todos los dispositivos envían un mensaje JSON como - { "ID" - "abc_001" } al tema - "tema/información".
¿Alguien puede sugerir cómo puedo superar este problema y manejarlo? ¿Puedo elegir la opción de tener varios procesos para ejecutar la función? Cualquier ayuda con respecto a lo mismo es bienvenida ya que estoy atascado en esto. Agradeciendotelo de antemano.
EDITAR: me gustaría aclarar que una sola instancia de este script se ejecutará en un servidor remoto. Habrá múltiples dispositivos enviando mensajes a este servidor remoto ejecutando este script usando MQTT. Por lo tanto, el suscriptor es este único cliente mqtt y hay varios editores.
const mqtt = require("mqtt"); const host = "0.0.0.0"; const port = "1883"; const connectURL = "mqtt://"+host + ":" + port; const subTopics = ["ABC", "XYZ","PQR"]; function subscribeTopics(client, ID){ topic = "tch/stb/"+ID.toString()+"/"; for (let i in subTopics) { client.subscribe(topic+subTopics[i], () => { console.log(`Subscribe to topic '${(topic+subTopics[i])}'`) }); } //Perform Database Insertion } const client = mqtt.connect(connectURL); client.on('connect', () => { console.log('Connected'); client.subscribe("topic/info", () => { console.log('Subscribe to topic topic/info') }); }) client.on('message', (topic, payload) => { message = JSON.parse(payload.toString()); if(topic === "topic/info"){ subscribeTopics(client, message.ID); } });Si tiene la intención de usar varias instancias del script para distribuir la carga, deberá usar un intermediario MQTT que admita suscripciones compartidas (esta es una función opcional agregada a la especificación en v5, pero algunos intermediarios tenían implementaciones propiciatorias en v3).
Sin él, todas las instancias del script recibirán TODOS los mensajes publicados en un tema para que todos procesen todos los mensajes. Con las suscripciones compartidas, cada mensaje solo se entregará a una instancia del cliente.
Esta publicación de blog debería ayudar a explicar las suscripciones compartidas con más detalle: https://www.hivemq.com/blog/mqtt5-essentials-part7-shared-subscriptions/