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

如何在Python中使用Async TCP实现服务器向PyQt客户端发送消息?

解决方案:实现PyQt+Async TCP聊天应用的服务器主动发消息

核心问题分析

你当前代码存在几个关键问题:

  • 单个global_writer只能维护一个客户端连接,多客户端场景下会失效
  • handle_client处理完一条消息就关闭连接,无法保持长连接进行双向通信
  • 未将GUI发送按钮触发的消息队列与服务器主动发送逻辑关联

修改后的完整实现代码

import asyncio
from PyQt5.QtWidgets import (QApplication, QWidget, QLineEdit, QPushButton, 
                             QComboBox, QVBoxLayout, QTextEdit)
import sys

class ChatApp(QWidget):
    def __init__(self):
        super().__init__()
        self._asyncio_event_loop = asyncio.new_event_loop()
        self._async_queue = asyncio.Queue()
        self.conType = ""
        self.init_ui()
        # 启动消息队列处理任务
        asyncio.run_coroutine_threadsafe(self._process_queue(), self._asyncio_event_loop)

    def init_ui(self):
        self.lineEdit = QLineEdit()
        self.sendButton = QPushButton("SEND")
        self.connectionButton = QPushButton("Connect")
        self.type1 = QComboBox()
        self.type1.addItems(["TCP Server", "TCP Client"])
        self.chatDisplay = QTextEdit()
        self.chatDisplay.setReadOnly(True)

        layout = QVBoxLayout()
        layout.addWidget(self.type1)
        layout.addWidget(self.connectionButton)
        layout.addWidget(self.lineEdit)
        layout.addWidget(self.sendButton)
        layout.addWidget(self.chatDisplay)
        self.setLayout(layout)

        self.connectionButton.pressed.connect(self.connection)
        self.sendButton.pressed.connect(lambda: self.send_message_to_event_loop(self.lineEdit.text(), self.conType))
        self.lineEdit.returnPressed.connect(lambda: self.send_message_to_event_loop(self.lineEdit.text(), self.conType))

    def connection(self):
        self.conType = self.type1.currentText()
        if self.conType == "TCP Server":
            self.send_message_to_event_loop("Server started", "TCP Server")
            # 启动TCP服务器
            asyncio.run_coroutine_threadsafe(self.start_server(), self._asyncio_event_loop)
        elif self.conType == "TCP Client":
            self.send_message_to_event_loop("Client started", "TCP Client")
            # 可在此补充客户端连接逻辑

    def send_message_to_event_loop(self, message: str, type_con) -> None:
        if not message:
            return
        asyncio.run_coroutine_threadsafe(
            self._async_queue.put((type_con, message)),
            loop=self._asyncio_event_loop,
        )
        # 显示自身发送的消息
        self.chatDisplay.append(f"[{type_con}] You: {message}")
        self.lineEdit.clear()

    async def _process_queue(self):
        """持续处理GUI传入事件循环的消息"""
        while True:
            msg_type, message = await self._async_queue.get()
            if msg_type == "TCP Server":
                # 服务器向所有活跃客户端广播消息
                await broadcast_message(message)
            elif msg_type == "TCP Client":
                # 客户端向服务器发送消息(补充客户端逻辑时实现)
                pass

# 服务器端活跃连接管理
active_writers = set()

async def broadcast_message(message: str):
    """向所有在线客户端广播消息"""
    message_bytes = f"Server: {message}".encode()
    # 遍历所有活跃连接发送消息,用list避免遍历中集合变化报错
    for writer in list(active_writers):
        try:
            writer.write(message_bytes)
            await writer.drain()
        except Exception as e:
            print(f"消息发送失败: {e}")
            active_writers.discard(writer)

async def handle_client(reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
    """处理单个客户端的连接与消息交互"""
    addr = writer.get_extra_info('peername')
    print(f"客户端 {addr} 已连接")
    active_writers.add(writer)

    try:
        while True:
            data = await reader.read(100)
            if not data:
                break
            message = data.decode().strip()
            print(f"从 {addr} 收到消息: '{message}'")
            # 将客户端消息同步到GUI
            app.window.send_message_to_event_loop(f"Client {addr}: {message}", "TCP Server")
            # 可选:将客户端消息广播给其他在线客户端
            await broadcast_message(f"Client {addr}: {message}")
    except Exception as e:
        print(f"客户端 {addr} 连接异常: {e}")
    finally:
        print(f"客户端 {addr} 已断开")
        active_writers.discard(writer)
        writer.close()
        await writer.wait_closed()

async def start_server():
    """启动TCP服务器"""
    server = await asyncio.start_server(handle_client, '127.0.0.1', 8888)
    addr = server.sockets[0].getsockname()
    print(f"服务器监听地址: {addr}")
    async with server:
        await server.serve_forever()

if __name__ == "__main__":
    app = QApplication(sys.argv)
    window = ChatApp()
    window.show()
    # 启动asyncio事件循环线程(后台运行)
    import threading
    threading.Thread(target=lambda: asyncio.set_event_loop(window._asyncio_event_loop) or window._asyncio_event_loop.run_forever(), daemon=True).start()
    sys.exit(app.exec_())

关键改动说明

  • 多客户端连接支持:用active_writers集合替代单个全局writer,自动维护所有在线客户端连接
  • 长连接保持:handle_client循环持续监听客户端消息,直到连接主动断开
  • 消息队列联动:_process_queue协程持续消费GUI传入的消息,触发服务器广播逻辑
  • 线程安全交互:所有跨线程操作通过asyncio.Queue完成,避免PyQt与asyncio的线程冲突
  • GUI反馈优化:发送消息后自动清空输入框并显示聊天记录,客户端消息同步展示

内容的提问来源于stack exchange,提问作者Barış Şenol

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 13:07:50