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

Python简易发布订阅应用中sock.recv()无响应问题求助

Python发布订阅程序:Broker广播消息客户端无法接收的解决方案

问题现象

实现了简易Pub/Sub程序,连接建立正常,客户端发送的消息Broker能接收,但Broker调用send()广播消息时,客户端的recv()无法收到内容。交互序列如下:

Broker                      Client 62863        Client 62867
------------------------------------------------------------
start()
                            start()
(62863) joined.
                                                start()
(62867) joined.
                            Hello from 62863
(62863): Hello from 62863
(62863) => (62867)

最后一步Broker执行广播,但Client 62867未收到消息。

问题根源

  1. Broker发送数据类型错误:
    Broker的listen_thread方法中,将接收到的字节流解码为字符串msg = publisher.recv(1024).decode(),但在broadcast方法中直接调用subscriber.send(msg)发送字符串。而socket的send()方法要求传入字节流(bytes),直接传字符串会导致消息无法正常发送。

  2. 客户端打印参数错误:
    Client的sock_to_stdout方法中,print(msg.decode('utf-8'), eol='', flush=True)使用了无效参数eol='',Python的print()函数没有该参数,正确参数是end=''。这个错误会导致线程抛出异常后终止,即使后续收到消息也无法处理。

修复方案

修改Broker代码

在broadcast方法中,将字符串编码为字节流后发送:

def broadcast(self, publisher, msg):
    for subscriber in self._subscribers:
        if publisher != subscriber:        # don't send to yourself
            print(f'{publisher.getpeername()} => {subscriber.getpeername()}', flush=True)
            try:
                subscriber.send(msg.encode('utf-8'))  # 编码为bytes后发送
            except:
                # broken socket, remove from subscriber list
                self._subscribers.remove(subscriber)

修改Client代码

修正print()函数的参数:

def sock_to_stdout(self):
    """
    Print anything received from the broker on stdout.
    """
    while True:
        msg = self._sock.recv(1024)
        print(msg.decode('utf-8'), end='', flush=True)  # 替换eol为end

完整修复后代码

完整Broker代码

import socket
import threading

class Broker(object):

    def __init__(self, host='', port=5000):
        self._host = host
        self._port = port
        self._subscribers = []

    def start(self):
        self._socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
        self._socket.bind((self._host, self._port))

        while True:
            """
            Wait for a client to request a connection and spawn a thread to
            receive and forward messages from that client.
            """
            self._socket.listen()
            subscriber, addr = self._socket.accept()
            print(f'{addr} joined.', flush=True)
            self._subscribers.append(subscriber)
            threading.Thread(target=self.listen_thread, args=(subscriber,)).start()

    def listen_thread(self, publisher):
        """
        Wait for a message to arrive from a publisher, broadcast to all other
        subscribers.
        """
        while True:
            msg = publisher.recv(1024).decode()
            if msg is not None:
                print(f'{publisher.getpeername()} published: {msg}', end='', flush=True)
                self.broadcast(publisher, msg)
            else:
                print(f'{publisher.getpeername()} has disconnected')
                return

    def broadcast(self, publisher, msg):
        for subscriber in self._subscribers:
            if publisher != subscriber:        # don't send to yourself
                print(f'{publisher.getpeername()} => {subscriber.getpeername()}', flush=True)
                try:
                    subscriber.send(msg.encode('utf-8'))
                except:
                    # broken socket, remove from subscriber list
                    self._subscribers.remove(subscriber)


if __name__ == "__main__":
    Broker().start()

完整Client代码

import socket
import threading
import sys

class StdioClient(object):
    """
    A simple pub/sub client:
    Anything received on stdin is published to the broker.
    Concurrently, anything broadcast by the broker is printed on stdout.
    """

    def __init__(self, host='localhost', port=5000):
        self._host = host
        self._port = port

    def start(self):
        self._sock = socket.socket()
        self._sock.connect((self._host, self._port))
        threading.Thread(target=self.stdin_to_sock).start()
        threading.Thread(target=self.sock_to_stdout).start()

    def stdin_to_sock(self):
        """
        Send anything received on stdin to the broker.
        """
        for msg in sys.stdin:
            self._sock.send(bytes(msg, 'utf-8'))

    def sock_to_stdout(self):
        """
        Print anything received from the broker on stdout.
        """
        while True:
            msg = self._sock.recv(1024)
            print(msg.decode('utf-8'), end='', flush=True)

if __name__ == '__main__':
    StdioClient().start()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 08:20:41