Python Asyncio中StreamReader读取与事件处理冲突问题求助
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
StreamReaderwithreadexactly()for sequential, reliable packet parsing - Use
asyncio.Protocolfor event-driven data handling with custom buffer management
内容的提问来源于stack exchange,提问作者Qwerty

