如何用asyncio.StreamReader.readline()处理超长行并继续读取?
asyncio 超长行安全读取替代方案
asyncio的readline()方法支持从流中异步读取单行,但当行长度超过限制时会抛出异常,且异常抛出后无法直接恢复读取。我们需要一个类似readline()的函数,能够自动丢弃行中超出限制的部分,并持续读取直到流结束。
现有方案情况
Python标准库中没有提供直接满足该需求的方法,需要自行实现。
实现思路
核心逻辑:
- 循环读取固定大小的数据块,累积到缓冲区
- 检查缓冲区中是否存在换行符:
- 若存在,提取换行符前的内容作为一行返回,剩余内容保留在缓冲区
- 若不存在,当缓冲区长度达到限制时,截断缓冲区到限制长度并返回,剩余内容继续处理
- 流结束时返回缓冲区中剩余的所有内容
代码实现
import asyncio async def safe_readline(reader: asyncio.StreamReader, limit: int = 2**16): buffer = bytearray() while True: # 每次读取最多limit字节的块 chunk = await reader.read(limit) if not chunk: # 流已结束,返回剩余内容(如果有) return buffer.decode() if buffer else None buffer.extend(chunk) newline_index = buffer.find(b'\n') if newline_index != -1: # 提取完整行(包含换行符) line = buffer[:newline_index + 1].decode() # 更新缓冲区为剩余内容 buffer = buffer[newline_index + 1:] return line elif len(buffer) >= limit: # 达到长度限制,截断并返回 line = buffer[:limit].decode() buffer = buffer[limit:] return line
使用示例
async def process_stream(reader: asyncio.StreamReader): while True: line = await safe_readline(reader) if line is None: break # 处理读取到的行 print(f"Received line: {line.rstrip()}") # 示例:创建连接并处理流 async def main(): reader, writer = await asyncio.open_connection('example.com', 80) writer.write(b'GET / HTTP/1.1\r\nHost: example.com\r\n\r\n') await writer.drain() await process_stream(reader) writer.close() await writer.wait_closed() asyncio.run(main())
关键说明
- 该函数与
readline()行为对齐,默认limit设为65536(2^16),和标准库一致 - 超长行处理:当缓冲区达到
limit时自动截断,剩余内容会在后续调用中继续处理 - 异常安全:无需处理
readline()抛出的超长行异常,全程稳定读取直到流结束
内容的提问来源于stack exchange,提问作者static_rtti
相关产品推荐
相关产品推荐

