PyZMQ客户端接收数据时无限挂起问题求助
ZeroMQ Pub-Sub模式客户端挂起问题修复
问题根源
- Pub-Sub消息丢失:服务器启动后立即发送所有消息,但此时客户端可能还未完成订阅连接。ZeroMQ Pub-Sub模式默认不保存历史消息,导致客户端收不到任何消息,持续阻塞在
subscriber.recv_json()。 - ROUTER Socket接收逻辑错误:ROUTER类型的Socket接收的消息包含客户端身份帧和实际消息两部分,直接调用
recv_json()会导致服务器阻塞,无法处理后续逻辑。
修复步骤
1. 调整服务器消息发送时机
在服务器发送消息前添加短暂延迟,确保客户端有足够时间完成订阅连接;或者改为等待客户端连接确认后再发送消息(更可靠)。
2. 修正ROUTER Socket接收逻辑
ROUTER接收消息时,需先接收客户端身份帧,再接收实际的JSON消息。
3. 优化清理逻辑
遍历消息字典副本进行删除操作,避免遍历过程中修改字典引发异常。
修改后的代码
Server.py
# -*- coding: utf-8 -*- import zmq import threading import json import os import time # 初始化上下文与Socket context = zmq.Context() publisher = context.socket(zmq.PUB) publisher.bind("tcp://*:5556") router = context.socket(zmq.ROUTER) router.bind("tcp://*:5557") # 模拟数据库存储消息状态 messages = {} consumers = ['consumer1', 'consumer2'] # 等待客户端完成订阅连接(延迟3秒,可根据实际情况调整) print("等待客户端连接...") time.sleep(3) # 发布10条消息 for i in range(10): message_id = str(i) file_path = f'./content/{i}' messages[message_id] = { 'file_path': file_path, 'consumers': consumers.copy(), 'processed_by': [], } publisher.send_json({ 'message_id': message_id, 'text': f'This is message {i}', 'media_path': file_path, }) print(f"已发送消息 {i}") # 清理已完成处理的消息 def cleanup(): while True: # 遍历字典副本,避免删除时引发异常 for message_id, message in list(messages.items()): if set(message['consumers']) == set(message['processed_by']): print(f"Deleting file {message['file_path']}") # os.remove(message['file_path']) # 取消注释可实际删除文件 del messages[message_id] time.sleep(5) cleanup_thread = threading.Thread(target=cleanup) cleanup_thread.start() # 接收客户端确认消息 while True: # ROUTER需先接收身份帧,再接收实际消息 identity = router.recv() ack_message = router.recv_json() print(f"收到来自 {identity.decode()} 的确认: {ack_message}") if ack_message['message_id'] in messages: consumer = ack_message['consumer'] # 避免重复添加确认记录 if consumer not in messages[ack_message['message_id']]['processed_by']: messages[ack_message['message_id']]['processed_by'].append(consumer)
Client.py
import zmq import time context = zmq.Context() consumer_id = 'consumer1' # 多客户端需修改此ID # 初始化订阅Socket subscriber = context.socket(zmq.SUB) subscriber.connect("tcp://localhost:5556") subscriber.setsockopt_string(zmq.SUBSCRIBE, '') # 初始化确认发送Socket dealer = context.socket(zmq.DEALER) dealer.identity = consumer_id.encode() dealer.connect("tcp://localhost:5557") # 处理消息循环 while True: message = subscriber.recv_json() print(f"收到消息: {message}") # 发送处理确认 dealer.send_json({ 'message_id': message['message_id'], 'consumer': consumer_id, }) time.sleep(1) # 测试阶段缩短延迟,方便观察
额外说明
- 如果需要持久化历史消息,可使用ZeroMQ的
XPUB/XSUB代理模式,或在服务器端维护消息队列,新客户端连接后主动推送未处理消息。 - Windows环境下需确保防火墙允许5556、5557端口的TCP连接。
内容的提问来源于stack exchange,提问作者sherwood
相关产品推荐
相关产品推荐

