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

如何判断asyncio中StreamReader何时就绪?TCP长连接场景

Handling Asyncio TCP Reader Readiness & Triggering Handlers

Great question! When you've got a persistent asyncio TCP connection and need to react automatically when the reader has incoming data, there are a few straightforward, idiomatic approaches in asyncio. Let's walk through them:

1. Persistent Coroutine for Reading (Most Common Approach)

The simplest way is to spin up a dedicated background coroutine that continuously waits for data using asyncio's built-in asynchronous read methods. These methods (like read(), readline(), or readexactly()) are designed to suspend the coroutine until data is available, so you don't have to manually check for readiness.

Here's a concrete example:

import asyncio

def process_data(data):
    # Replace this with your actual business logic
    print(f"Processing received data: {data.decode().strip()}")

async def handle_incoming(reader):
    while True:
        # Adjust the read method based on your protocol (e.g., readline() for line-based data)
        data = await reader.read(1024)
        if not data:
            # Empty bytes mean the remote end closed the connection
            print("Connection closed by peer")
            break
        # Trigger your data handler
        process_data(data)

async def main():
    addr = ("localhost", 8888)
    reader, writer = await asyncio.open_connection(*addr)
    
    # Start the data handling coroutine in the background
    asyncio.create_task(handle_incoming(reader))
    
    # Keep the main coroutine alive (replace with your app's main logic if needed)
    await asyncio.Future()

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

This approach fits naturally with asyncio's asynchronous model, requires no low-level event loop fiddling, and keeps your code clean.

2. Low-Level Event Loop add_reader (Fine-Grained Control)

If you need more direct control over IO events, you can use the event loop's add_reader() method to listen for the reader's file descriptor becoming readable. When data is available, a synchronous callback is triggered.

Note: The callback runs synchronously, so avoid blocking operations here—use asyncio.create_task() if you need to run async logic inside it.

Example:

import asyncio

def on_data_ready(reader):
    # Use read_nowait() since we know data is available (triggered by add_reader)
    try:
        data = reader.read_nowait()
        if data:
            print(f"Received via direct FD listen: {data.decode().strip()}")
            process_data(data)
        else:
            # Connection closed, remove the listener
            loop = asyncio.get_running_loop()
            loop.remove_reader(reader.fileno())
            print("Connection closed")
    except asyncio.IncompleteReadError:
        # Handle partial reads if needed
        pass

async def main():
    addr = ("localhost", 8888)
    reader, writer = await asyncio.open_connection(*addr)
    
    loop = asyncio.get_running_loop()
    # Register the listener for the reader's file descriptor
    loop.add_reader(reader.fileno(), on_data_ready, reader)
    
    await asyncio.Future()

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

This is useful for scenarios where you need to integrate with non-asyncio IO code or want precise control over event handling.

3. Decouple Reading & Processing with a Queue (Complex Workflows)

For more complex applications where data processing needs to be separated from IO reading (e.g., batching, multiple processors), use an asyncio.Queue to pass data between the reader coroutine and your processing logic.

Example:

import asyncio

async def reader_coroutine(reader, queue):
    while True:
        data = await reader.read(1024)
        if not data:
            # Send a sentinel value to signal connection closure
            await queue.put(None)
            break
        await queue.put(data)

async def processor_coroutine(queue):
    while True:
        data = await queue.get()
        if data is None:
            print("Processor shutting down: connection closed")
            break
        print(f"Processor handling data: {data.decode().strip()}")
        process_data(data)
        queue.task_done()

async def main():
    addr = ("localhost", 8888)
    reader, writer = await asyncio.open_connection(*addr)
    
    data_queue = asyncio.Queue()
    # Start both reader and processor tasks
    asyncio.create_task(reader_coroutine(reader, data_queue))
    asyncio.create_task(processor_coroutine(data_queue))
    
    await asyncio.Future()

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

This pattern makes your code more modular and easier to scale as your application grows.

Key Notes

  • Always handle the case where read() returns empty bytes—this indicates the remote connection has been closed.
  • When sending data, use writer.write(data) followed by await writer.drain() to ensure data is flushed properly.
  • Don't forget to clean up resources: call writer.close() and await writer.wait_closed() when your connection is no longer needed.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 04:21:27