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))
问题根源
- 错误使用XPUB套接字:XPUB是代理模式专用的发布套接字,需要处理订阅者的订阅信号,普通双向通信场景应该用
zmq.PUB/zmq.SUB或zmq.REQ/zmq.REP,XPUB默认会在无订阅者时阻塞发送。 - XML数据格式错误:
et.tostring()返回字节串,直接和字符串拼接后发送,会导致订阅端解析失败。 - 无连接等待机制: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
相关产品推荐
相关产品推荐

