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

Python中ZMQ同时收发消息阻塞问题求助

ZMQ 双向通信阻塞问题排查与解决

我在学习Python的ZMQ库,想实现Publisher发送消息后接收Subscriber的确认消息,但现在遇到阻塞问题:发布端卡在"Data Sent",订阅端卡在"Getting Message...",无法正常收发。


原代码

发布端(Publisher)

import xml.etree.ElementTree as et
import zmq
import os
import time
from datetime import datetime

# Publish connections
outPort = int(os.environ.get('OUTPORT', "8760"))  # GMH 0MQ bus
contextPub = zmq.Context()
socketPub = contextPub.socket(zmq.XPUB)
socketPub.bind("tcp://*:%s" % outPort)
outTopic = "xmltest"

# Subscription connections
dataSource = os.environ.get('DATASOURCE', 'tcp://localhost:8761')
contextSub = zmq.Context()
socketSub = contextSub.socket(zmq.SUB)
socketSub.connect(dataSource)
topicfilter = "c2test"
socketSub.setsockopt_string(zmq.SUBSCRIBE, topicfilter)

# Create and test writing xml
data = et.Element('myMessage')

element1 = et.SubElement(data, 'simpleType')

sElem1 = et.SubElement(element1, 'Freq_Center_MHz')

sElem2 = et.SubElement(element1, 'Bandwidth_MHz')

sElem3 = et.SubElement(element1, 'samplerate_msps')

sElem4 = et.SubElement(element1, 'time')

sElem1.text = '100'
sElem2.text = '56.1'
sElem3.text = '56.1'

while True:
    myDateTime = datetime.now().strftime("%Y%m%dT%H:%M:%S")
    sElem4.text = str(myDateTime)
    bXml = et.tostring(data)
    socketPub.send_string("%s %s" % (outTopic, bXml))
    print("Data Sent")
    time.sleep(1)

    messageStr = socketSub.recv(flags=0)
    print("Hello")
    print(messageStr)

socketPub.close()

订阅端(Subscriber)

import xml.etree.ElementTree as et
import zmq
import os

# Publish connections
outPort = int(os.environ.get('OUTPORT', "8761"))  # GMH 0MQ bus
contextPub = zmq.Context()
socketPub = contextPub.socket(zmq.XPUB)
socketPub.bind("tcp://*:%s" % outPort)
outTopic = "c2test"

# Subscription connections
dataSource = os.environ.get('DATASOURCE', 'tcp://localhost:8760')
contextSub = zmq.Context()
socketSub = contextSub.socket(zmq.SUB)
socketSub.connect(dataSource)
topicfilter = "xmltest"
socketSub.setsockopt_string(zmq.SUBSCRIBE, topicfilter)

while True:
    print("Getting Message...")
    messageStr = socketSub.recv()
    print("Message Received")
    res = messageStr.decode().split(" ")[1].split("'")[1]

    tree = et.fromstring(res)
    for child in tree.findall("./simpleType/*"):
        print("%s: %s" % (child.tag, child.text))

    print()
    print("Sending Acknowledgement.")
    print()
    receive = "Message Received"
    socketPub.send_string("%s %s" % (outTopic, receive))

问题根源

  1. 错误使用XPUB套接字:XPUB是代理模式专用的发布套接字,需要处理订阅者的订阅信号,普通双向通信场景应该用zmq.PUB/zmq.SUB或zmq.REQ/zmq.REP,XPUB默认会在无订阅者时阻塞发送。
  2. XML数据格式错误:et.tostring()返回字节串,直接和字符串拼接后发送,会导致订阅端解析失败。
  3. 无连接等待机制:PUB套接字需要等待订阅者建立连接才能发送消息,原代码没有预留连接时间,导致消息无法送达。

修复方案

方案1:REQ-REP模式(适合一对一请求-确认场景)

发布端(REQ)

import xml.etree.ElementTree as et
import zmq
import time
from datetime import datetime

context = zmq.Context()
socket = context.socket(zmq.REQ)
socket.connect("tcp://localhost:8760")

# 构建XML结构
data = et.Element('myMessage')
element1 = et.SubElement(data, 'simpleType')
et.SubElement(element1, 'Freq_Center_MHz').text = '100'
et.SubElement(element1, 'Bandwidth_MHz').text = '56.1'
et.SubElement(element1, 'samplerate_msps').text = '56.1'
time_elem = et.SubElement(element1, 'time')

while True:
    # 更新时间字段
    time_elem.text = datetime.now().strftime("%Y%m%dT%H:%M:%S")
    # 发送XML字节数据
    socket.send(et.tostring(data))
    print("Data Sent")
    
    # 接收确认消息
    ack = socket.recv_string()
    print(f"Received Acknowledgement: {ack}")
    time.sleep(1)

socket.close()
context.term()

订阅端(REP)

import xml.etree.ElementTree as et
import zmq

context = zmq.Context()
socket = context.socket(zmq.REP)
socket.bind("tcp://*:8760")

while True:
    print("Getting Message...")
    # 接收XML字节数据
    xml_bytes = socket.recv()
    print("Message Received")
    
    # 解析并打印XML内容
    tree = et.fromstring(xml_bytes)
    for child in tree.findall("./simpleType/*"):
        print(f"{child.tag}: {child.text}")
    
    # 发送确认消息
    socket.send_string("Message Received")
    print("\nSent Acknowledgement.\n")

socket.close()
context.term()

方案2:PUB-SUB双向主题模式(适合一对多场景)

发布端

import xml.etree.ElementTree as et
import zmq
import time
from datetime import datetime

# 发布消息用PUB
pub_context = zmq.Context()
pub_socket = pub_context.socket(zmq.PUB)
pub_socket.bind("tcp://*:8760")
out_topic = "xmltest"

# 接收确认用SUB
sub_context = zmq.Context()
sub_socket = sub_context.socket(zmq.SUB)
sub_socket.connect("tcp://localhost:8761")
sub_socket.setsockopt_string(zmq.SUBSCRIBE, "c2test")

# 构建XML结构
data = et.Element('myMessage')
element1 = et.SubElement(data, 'simpleType')
et.SubElement(element1, 'Freq_Center_MHz').text = '100'
et.SubElement(element1, 'Bandwidth_MHz').text = '56.1'
et.SubElement(element1, 'samplerate_msps').text = '56.1'
time_elem = et.SubElement(element1, 'time')

# 等待订阅者连接
time.sleep(2)

while True:
    time_elem.text = datetime.now().strftime("%Y%m%dT%H:%M:%S")
    # 将XML转为字符串发送
    xml_str = et.tostring(data, encoding='unicode')
    pub_socket.send_string(f"{out_topic} {xml_str}")
    print("Data Sent")
    
    # 非阻塞接收确认,避免无限阻塞
    if sub_socket.poll(1000) != 0:
        message = sub_socket.recv_string()
        print(f"Received: {message}")
    else:
        print("No acknowledgement received within timeout")
    
    time.sleep(1)

pub_socket.close()
sub_socket.close()
pub_context.term()
sub_context.term()

订阅端

import xml.etree.ElementTree as et
import zmq

# 订阅消息用SUB
sub_context = zmq.Context()
sub_socket = sub_context.socket(zmq.SUB)
sub_socket.connect("tcp://localhost:8760")
sub_socket.setsockopt_string(zmq.SUBSCRIBE, "xmltest")

# 发送确认用PUB
pub_context = zmq.Context()
pub_socket = pub_context.socket(zmq.PUB)
pub_socket.bind("tcp://*:8761")
out_topic = "c2test"

while True:
    print("Getting Message...")
    # 拆分主题和消息内容
    message = sub_socket.recv_string()
    topic, xml_str = message.split(" ", 1)
    print("Message Received")
    
    # 解析XML
    tree = et.fromstring(xml_str)
    for child in tree.findall("./simpleType/*"):
        print(f"{child.tag}: {child.text}")
    
    # 发送确认
    pub_socket.send_string(f"{out_topic} Message Received")
    print("\nSent Acknowledgement.\n")

sub_socket.close()
pub_socket.close()
sub_context.term()
pub_context.term()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 13:20:30