You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于Python实现ZeroMQ网络的实时图形可视化?

实现ZeroMQ Pub/Sub网络自动可视化的方案

1. 基于ZeroMQ监控API采集连接数据

ZeroMQ内置的ZMQ_MONITOR是实现自动发现的核心,它能捕获套接字的连接、断开等事件,无需大幅修改原有节点业务逻辑。给每个节点的Pub/Sub套接字添加监控线程即可:

import zmq
import threading

def monitor_socket(socket, visualizer_endpoint):
    # 创建监控套接字,订阅所有事件
    monitor = socket.get_monitor_socket()
    monitor.subscribe(b"")
    
    # 连接到可视化中心节点
    ctx = zmq.Context.instance()
    reporter = ctx.socket(zmq.PUSH)
    reporter.connect(visualizer_endpoint)

    while True:
        try:
            event_parts = monitor.recv_multipart()
            event_type = event_parts[0].decode()
            local_addr = event_parts[2].decode()
            remote_addr = event_parts[3].decode()

            # 只处理有效连接事件
            if event_type in ["CONNECTED", "ACCEPTED"]:
                # 发送连接信息到可视化中心
                reporter.send_json({
                    "local": local_addr,
                    "remote": remote_addr,
                    "event": event_type
                })
        except zmq.ZMQError as e:
            if e.errno != zmq.EAGAIN:
                break

# 原有节点示例:创建Pub套接字并启动监控
ctx = zmq.Context()
pub_socket = ctx.socket(zmq.PUB)
pub_socket.bind("tcp://*:5555")
# 启动监控线程,指定可视化中心地址
threading.Thread(target=monitor_socket, args=(pub_socket, "tcp://localhost:5556"), daemon=True).start()

2. 搭建拓扑收集中心

单独编写一个可视化节点,负责接收所有运行节点的监控数据,维护全局拓扑结构:

import zmq
import time
from collections import defaultdict

class TopologyCollector:
    def __init__(self):
        self.topology = defaultdict(list)
        self.ctx = zmq.Context()
        self.receiver = self.ctx.socket(zmq.PULL)
        self.receiver.bind("tcp://*:5556")

    def update_topology(self):
        while True:
            try:
                msg = self.receiver.recv_json(flags=zmq.NOBLOCK)
                local = msg["local"]
                remote = msg["remote"]
                # 区分Pub(ACCEPTED事件的本地端)和Sub(CONNECTED事件的本地端)
                if msg["event"] == "ACCEPTED":
                    pub_node = local
                    sub_node = remote
                else:
                    pub_node = remote
                    sub_node = local
                # 去重添加连接
                if sub_node not in self.topology[pub_node]:
                    self.topology[pub_node].append(sub_node)
            except zmq.ZMQError as e:
                if e.errno == zmq.EAGAIN:
                    time.sleep(0.5)
                    continue
                break

3. 图形化渲染拓扑

用networkx+matplotlib生成静态可视化,或用pyvis生成可交互的动态页面:

静态渲染示例(networkx + matplotlib)

import networkx as nx
import matplotlib.pyplot as plt
from topology_collector import TopologyCollector

def draw_topology(collector):
    G = nx.DiGraph()
    # 添加节点和边
    for pub, subs in collector.topology.items():
        G.add_node(pub, label=f"Pub\n{pub.split(':')[-1]}", color="#2ecc71")
        for sub in subs:
            G.add_node(sub, label=f"Sub\n{sub.split(':')[-1]}", color="#3498db")
            G.add_edge(pub, sub)
    # 布局与绘制
    pos = nx.spring_layout(G, k=0.5)
    node_colors = [G.nodes[n]["color"] for n in G.nodes]
    nx.draw(G, pos, with_labels=True, node_color=node_colors, node_size=2500, font_size=8)
    plt.title("ZeroMQ Pub/Sub Network Topology")
    plt.show()

# 运行流程
collector = TopologyCollector()
# 启动收集线程
import threading
threading.Thread(target=collector.update_topology, daemon=True).start()
# 定期刷新绘图
while True:
    plt.clf()
    draw_topology(collector)
    plt.pause(2)

4. 动态节点发现补充(Beacon机制)

针对动态启动的节点,用ZeroMQ Beacon实现自动上报:

  • 每个节点启动Beacon,定期广播自身端点信息
  • 可视化中心监听Beacon广播,自动发现新节点并同步监控数据

关键注意事项

  • 合并重复事件:CONNECTED(Sub端发起)和ACCEPTED(Pub端接收)对应同一连接,需统一处理避免重复记录
  • 端点标识优化:尽量用带节点名称的地址(如tcp://order-service:5555),而非tcp://*:5555,提升可视化可读性
  • 失效连接清理:监听DISCONNECTED事件,及时从拓扑中移除失效连接

内容的提问来源于stack exchange,提问作者hunterlineage

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 17:02:23