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

Python Asyncio中StreamReader读取与事件处理冲突问题求助

Fixing StreamReader Conflict Between yield from read() and Custom feed_data

Hey there! Let's break down why you're hitting this issue and walk through the fixes.

The Root of the Problem

First off, StreamReader.feed_data is an internal method that asyncio uses to push incoming data into the reader's internal buffer. When you overwrite it with self.sock_reader.feed_data = self.on_data, you're breaking the core functionality of StreamReader—the read(1024) call relies on that buffer to fetch data, so your yield from statement will never get meaningful data once you've replaced feed_data.

Solution 1: Use StreamReader Properly with Packet Handling

Instead of messing with feed_data, handle your packet parsing directly after reading from the stream. Since your packets have a fixed 8-byte header (with code and length), use readexactly() to ensure you get complete chunks of data—this also solves "sticky packet" issues automatically.

Here's adjusted code for this approach:

import asyncio
import struct
import json

class YourStreamHandler:
    def __init__(self, sock_reader):
        self.sock_reader = sock_reader

    def on_data(self, code, payload):
        # Update your handler to work with parsed values instead of raw data
        print(f"Received code: {code}, payload: {payload}")

    async def process_incoming_data(self):
        while True:
            try:
                # Step 1: Read the full 8-byte header
                header = await self.sock_reader.readexactly(8)
                code, payload_length = struct.unpack('<ii', header)
                
                # Step 2: Read the exact payload length specified in the header
                payload_bytes = await self.sock_reader.readexactly(payload_length)
                payload = json.loads(payload_bytes.decode('utf-8'))
                
                # Step 3: Pass parsed data to your handler
                self.on_data(code, payload)
            except asyncio.IncompleteReadError:
                # Triggered when the connection closes unexpectedly
                print("Connection closed by peer")
                break

To use this handler:

async def main():
    # Set up your StreamReader (example via asyncio.open_connection)
    reader, writer = await asyncio.open_connection('localhost', 8888)
    handler = YourStreamHandler(reader)
    await handler.process_incoming_data()

asyncio.run(main())

Solution 2: Use asyncio.Protocol for Event-Driven Control

If you prefer an event-driven flow (similar to what you tried with feed_data), use asyncio.Protocol instead of StreamReader. This lets you handle incoming data directly as it arrives, while managing your own buffer to parse complete packets.

import asyncio
import struct
import json

class PacketProtocol(asyncio.Protocol):
    def __init__(self):
        self.buffer = b''  # Buffer to hold incomplete packets

    def connection_made(self, transport):
        self.transport = transport
        print("Connected to server")

    def data_received(self, data):
        self.buffer += data
        
        # Process all complete packets in the buffer
        while len(self.buffer) >= 8:
            # Check if we have enough data for the full packet (header + payload)
            header = self.buffer[:8]
            code, payload_length = struct.unpack('<ii', header)
            total_packet_size = 8 + payload_length
            
            if len(self.buffer) >= total_packet_size:
                # Extract the full packet and trim the buffer
                packet = self.buffer[:total_packet_size]
                self.buffer = self.buffer[total_packet_size:]
                
                # Parse and handle the packet
                self.on_data(packet)
            else:
                # Not enough data for a full packet—wait for more incoming bytes
                break

    def on_data(self, data):
        code, length = struct.unpack('<ii', data[:8])
        payload = json.loads(data[8:].decode('utf-8'))
        print(f"Processed packet: code={code}, payload={payload}")

    def connection_lost(self, exc):
        print("Connection lost")

To use this protocol:

async def main():
    loop = asyncio.get_running_loop()
    transport, protocol = await loop.create_connection(
        lambda: PacketProtocol(),
        'localhost', 8888
    )
    await asyncio.Future()  # Keep the connection alive indefinitely

asyncio.run(main())

Key Takeaway

You can't mix StreamReader's built-in read methods with overwriting feed_data—it breaks the reader's internal buffer logic. Pick one approach:

  • Use StreamReader with readexactly() for sequential, reliable packet parsing
  • Use asyncio.Protocol for event-driven data handling with custom buffer management

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.26 10:14:09