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

如何修改Python asyncio Socket服务器实现全客户端消息同步推送?

解决asyncio服务器同步推送相同消息给所有客户端的方案

核心思路

  • 维护一个线程安全的客户端连接集合,用asyncio.Lock保护并发修改
  • 单独启动异步任务负责定时生成随机消息
  • 消息生成后遍历所有在线客户端,统一推送相同内容

完整服务器代码

import asyncio
import random
import time

# 存储所有客户端的writer对象,加锁保护并发访问
clients = set()
clients_lock = asyncio.Lock()

async def generate_and_broadcast():
    """异步任务:定时生成全局消息并推送给所有客户端"""
    while True:
        # 随机等待1-5秒模拟随机触发时刻
        await asyncio.sleep(random.randint(1, 5))
        # 生成统一的消息内容
        message = f"[{time.strftime('%H:%M:%S')}] 全局推送消息: {random.randint(1000, 9999)}"
        message_bytes = (message + "\n").encode('utf-8')
        
        # 加锁复制当前客户端列表,避免遍历过程中集合被修改
        async with clients_lock:
            current_clients = list(clients)
        
        # 逐个推送给所有客户端
        for writer in current_clients:
            try:
                writer.write(message_bytes)
                await writer.drain()
            except Exception as e:
                # 推送失败,移除失效客户端
                print(f"推送失败,移除客户端: {e}")
                async with clients_lock:
                    if writer in clients:
                        clients.remove(writer)

async def handle_client(reader, writer):
    """处理单个客户端的连接生命周期"""
    # 新客户端加入连接池
    async with clients_lock:
        clients.add(writer)
        print(f"新客户端连接,当前在线: {len(clients)}")
    
    try:
        # 保持连接,若需处理客户端输入可在此添加逻辑
        while True:
            data = await reader.read(100)
            if not data:
                break
    finally:
        # 客户端断开,从连接池移除
        async with clients_lock:
            if writer in clients:
                clients.remove(writer)
            print(f"客户端断开,当前在线: {len(clients)}")
        writer.close()
        await writer.wait_closed()

async def main():
    # 启动TCP服务器
    server = await asyncio.start_server(handle_client, '127.0.0.1', 8888)
    print(f"服务器启动,监听 {server.sockets[0].getsockname()}")
    
    # 启动全局广播任务
    asyncio.create_task(generate_and_broadcast())
    
    # 持续运行服务器
    async with server:
        await server.serve_forever()

if __name__ == "__main__":
    asyncio.run(main())

客户端测试代码

import socket

def client():
    s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    s.connect(('127.0.0.1', 8888))
    print("连接服务器成功,等待消息...")
    while True:
        data = s.recv(1024)
        if not data:
            print("服务器断开连接")
            break
        print(data.decode('utf-8').strip())
    s.close()

if __name__ == "__main__":
    client()

关键细节说明

  • 连接池安全:用asyncio.Lock保护客户端集合的增删操作,避免并发修改导致的异常
  • 独立广播任务:消息生成与推送逻辑单独抽离,避免和客户端连接任务互相干扰,确保所有客户端收到同一条消息
  • 异常处理:推送失败时自动移除失效客户端,避免后续推送持续报错
  • 集合复制:遍历前复制客户端列表,防止遍历过程中集合被修改(如客户端断开)引发遍历异常

解决你之前的问题

  1. 针对"asynchronous generator is already running"错误:现在用独立的异步任务生成消息,客户端连接任务仅负责维护连接,不会复用同一个生成器,彻底解决该问题
  2. 针对asyncio.all_tasks()的使用:无需依赖系统任务列表,显式维护客户端连接池更清晰可控,避免误操作其他系统任务

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 16:41:01