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

Python中StreamReader.readline()读取换行符后未返回的问题

问题:StreamReader.readline()未按预期返回,需客户端关闭才接收数据

启动服务端后输出 Serving on ('127.0.0.1', 7777),运行客户端后服务端能收到初始消息 Received 'Service: TEST' from ('127.0.0.1', 39923),但后续客户端输入的内容需等客户端关闭连接才会被服务端打印。已确认客户端发送的数据包含换行符,但服务端的StreamReader.readline()未在收到换行符时立即返回。

相关代码

ServiceSubscription.py

import asyncio
from asyncio import StreamReader, StreamWriter

# 补充原代码缺失的Connections类
class Connections:
    def __init__(self, reader: StreamReader, writer: StreamWriter, msg: str):
        self.reader = reader
        self.writer = writer
        self.msg = msg

class ServiceSubscription():
    def __init__(self) -> None:
        self.subscriber_connections = []
        self.service_connections = []
        self.server_listener = None
        # 键为服务名称,值为对应ServiceReader实例
        self.service_readers = {}
    
    """
    在7777端口启动监听服务器
    """
    async def initiate_server(self):
        server = await asyncio.start_server(self.handle_incoming, '127.0.0.1', 7777)
        addrs = ', '.join(str(sock.getsockname()) for sock in server.sockets)
        print(f'Serving on {addrs}')
        
        async with server:
            await server.serve_forever()
    
    """
    根据连接类型(服务端/订阅端)处理入站连接
    首次发送的消息格式应为:'service:SERVICE_NAME' 或 'suscriber: [SERVICE1, SERVICE2, ...]'
    """
    async def handle_incoming(self, reader: StreamReader, writer: StreamWriter):
        # 改用readline读取初始命令,确保完整读取一行
        data = await reader.readline()
        message = data.decode().strip()
        addr = writer.get_extra_info('peername')

        print(f"Received {message!r} from {addr!r}")
        if "Service:" in message:
            # 正确提取服务名称
            service_name = message.split(':')[1].strip()
            self.service_connections.append(Connections(reader, writer, service_name))
            service_reader = ServiceReader(reader=reader, writer=writer)
            self.service_readers[service_name] = service_reader
            await service_reader.broadcast()
        
        elif "Suscriber:" in message:
            # 正确提取订阅的服务列表
            sub_services = message.split(':')[1].strip().strip('[]').split(',')
            sub_services = [s.strip() for s in sub_services]
            self.subscriber_connections.append(Connections(reader, writer, ','.join(sub_services)))
            # 为每个订阅的服务添加订阅者
            for service in sub_services:
                if service in self.service_readers:
                    self.service_readers[service].add_suscribers(writer)
        
        else:
            pass

class ServiceReader():
    def __init__(self, reader:  StreamReader, writer: StreamWriter):
        self.reader = reader
        self.writer = writer
        self.suscribers: list[StreamWriter] = []
        self._stop = asyncio.Event()

    def stop(self):
        self._stop.set()
        
    """
    添加新订阅者的StreamWriter实例
    """
    def add_suscribers(self, writer: StreamWriter):
        self.suscribers.append(writer)
        
    """
    读取服务端数据并广播给所有订阅者
    """
    async def broadcast(self):
        # 同时监听停止信号和数据读取
        while not self._stop.is_set():
            try:
                # 超时处理避免无限等待
                data = await asyncio.wait_for(self.reader.readline(), timeout=1)
                if not data:
                    break  # 连接关闭
                if b'\n' in data:
                    print(True)
                data_str = data.decode().strip()
                print(f"Received from service: {data_str}")
                # 广播数据给所有订阅者
                for subscriber in self.suscribers:
                    subscriber.write(data)
                    await subscriber.drain()
            except asyncio.TimeoutError:
                continue
        # 关闭当前连接
        self.writer.close()
        await self.writer.wait_closed()

WriterTest.py

import asyncio

async def async_input(prompt: str):
    # 使用线程池执行同步input,避免阻塞事件循环
    loop = asyncio.get_event_loop()
    return await loop.run_in_executor(None, input, prompt)

async def tcp_echo_client(message):
    reader, writer = await asyncio.open_connection(
        '127.0.0.1', 7777)

    print(f'Send: {message!r}\n')
    # 初始消息添加换行符,让服务端能准确读取完整命令
    writer.write((message + "\n").encode())
    await writer.drain()
    
    try:
        while True:
            data = await async_input("Type a message\n")
            if data.lower() == 'exit':
                break
            data = (data + "\n").encode()
            writer.write(data)
            await writer.drain()
    finally:
        writer.close()
        await writer.wait_closed()
    

asyncio.run(tcp_echo_client('Service: TEST'))

问题原因及修复方案

1. 客户端同步input阻塞事件循环

原客户端代码使用input()同步阻塞函数,会卡住asyncio事件循环,导致数据无法及时发送到服务端,服务端readline()只能等到客户端关闭时才收到所有缓冲数据。
修复:用asyncio.run_in_executor在单独线程执行input(),避免阻塞事件循环,确保数据实时发送。

2. 服务端初始读取方式错误

原服务端用reader.read(100)读取初始消息,无法保证读取完整命令,还可能缓冲后续数据干扰readline()工作。
修复:客户端初始消息添加换行符,服务端改用readline()读取初始命令,同时正确提取服务名称和订阅列表。

3. ServiceReader循环逻辑不完善

原broadcast()仅依赖reader.at_eof(),无停止信号和广播逻辑。修复后添加停止事件、超时处理,同时实现向订阅者广播数据的功能。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 09:31:14