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

PyZMQ客户端接收数据时无限挂起问题求助

ZeroMQ Pub-Sub模式客户端挂起问题修复

问题根源

  1. Pub-Sub消息丢失:服务器启动后立即发送所有消息,但此时客户端可能还未完成订阅连接。ZeroMQ Pub-Sub模式默认不保存历史消息,导致客户端收不到任何消息,持续阻塞在subscriber.recv_json()。
  2. 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 05:44:55