如何在Angular(ngx-mqtt)中建立多个MQTT连接并订阅各连接对应主题
Got it, let's figure out how to set up multiple MQTT connections and handle per-connection topic subscriptions with your existing code setup. Here's a practical, step-by-step approach tailored to your needs:
Right now, you're injecting a single MqttService instance, which won't cut it for multiple independent connections. We'll create a connection manager to wrap each connection's details and logic, so you can track and interact with them separately.
First, define a structure to hold all critical info for each connection—this makes it easy to manage subscriptions, client IDs, and connection status:
// Define an interface to encapsulate each MQTT connection export interface MqttConnection { hostname: string; clientId: string; port: number; client: any; // Replace with your actual MQTT client type (e.g., MqttClient) subscribedTopics: string[]; }
This service will handle creating new connections, subscribing to topics per connection, and cleaning up connections when needed. It uses Angular's Injector to dynamically create new MqttService instances (assuming each instance maps to one connection):
import { Injectable, Injector } from '@angular/core'; import { MqttService } from 'path-to-your-mqtt-service'; import { MqttConnection } from './mqtt-connection.interface'; @Injectable({ providedIn: 'root' }) export class MqttConnectionManager { private activeConnections: MqttConnection[] = []; constructor(private injector: Injector) {} // Create a new MQTT connection using your array3 values createConnection(hostname: string, clientId: string, port: number): MqttConnection { // Dynamically inject a new MqttService instance for this connection const mqttService = this.injector.get(MqttService); // Update your connect method to accept full connection params (see note below) const client = mqttService.connect({ hostname, clientId, port }); const newConnection: MqttConnection = { hostname, clientId, port, client, subscribedTopics: [] }; this.activeConnections.push(newConnection); return newConnection; } // Subscribe a specific connection to a topic subscribeToTopic(connection: MqttConnection, topic: string): void { if (!connection.subscribedTopics.includes(topic)) { connection.client.subscribe(topic); connection.subscribedTopics.push(topic); // Add message handling for this topic/connection connection.client.on('message', (receivedTopic: string, payload: Buffer) => { if (receivedTopic === topic) { console.log(`Connection ${connection.clientId} received message on ${topic}:`, payload.toString()); // Replace this with your custom message handling logic (e.g., emit to components) } }); } } // Get all active connections getAllActiveConnections(): MqttConnection[] { return [...this.activeConnections]; } // Disconnect and clean up a connection disconnectConnection(connection: MqttConnection): void { connection.client.end(); this.activeConnections = this.activeConnections.filter(conn => conn !== connection); } }
Your original connect call only uses array3[0] (hostname)—you'll need to modify it to accept all connection parameters:
// Inside your MqttService connect(options: { hostname: string; clientId: string; port: number }): any { // Adjust the connection string based on your broker's protocol (mqtt://, mqtts://, etc.) const connectionUrl = `mqtt://${options.hostname}:${options.port}`; // Create and return the client instance with the unique client ID return mqtt.connect(connectionUrl, { clientId: options.clientId }); }
Now you can loop through your array3 to create multiple connections and subscribe each to their respective topics:
import { Component, OnInit, OnDestroy } from '@angular/core'; import { MqttConnectionManager } from './mqtt-connection-manager.service'; @Component({ selector: 'app-mqtt-handler', templateUrl: './mqtt-handler.component.html' }) export class MqttHandlerComponent implements OnInit, OnDestroy { // Your array3: each entry is [hostname, clientId, port] array3: [string, string, number][] = [ ['mqtt-broker-1.com', 'client-device-001', 1883], ['mqtt-broker-2.com', 'client-device-002', 8883] ]; constructor(private connectionManager: MqttConnectionManager) {} ngOnInit(): void { // Create a connection for each entry in array3 this.array3.forEach(([hostname, clientId, port]) => { const newConnection = this.connectionManager.createConnection(hostname, clientId, port); // Subscribe to a topic specific to this connection (adjust the topic as needed) this.connectionManager.subscribeToTopic(newConnection, `device/${clientId}/status`); }); } ngOnDestroy(): void { // Clean up all connections when the component is destroyed this.connectionManager.getAllActiveConnections().forEach(conn => { this.connectionManager.disconnectConnection(conn); }); } }
- Unique Client IDs: Every connection needs a distinct
clientId—MQTT brokers will automatically drop old connections if a duplicate ID connects. Double-check yourarray3values! - Error Handling: Add logic to handle connection failures (e.g.,
client.on('error', (err) => { ... })) in the connection manager. - Protocol Adjustments: If your broker uses SSL/TLS, replace
mqtt://withmqtts://in the connection URL.
内容的提问来源于stack exchange,提问作者Ionic_pramod

