如何实现兼容同步/异步read方法的通用数据帧解析器?
通用可变长度数据帧解析器的同步/异步兼容方案
核心思路:分离IO操作与解析逻辑
不管是同步还是异步场景,数据帧的核心解析规则(格式校验、数据转换、校验和验证)是完全一致的,差异仅在于「读取指定字节数」这个IO操作的调用方式。因此我们只需要把核心解析逻辑抽离,让IO操作作为可替换的依赖传入,就能避免代码重复。
方案一:生成器驱动的核心复用
用生成器表达解析步骤(每次yield需要读取的字节数),再分别编写同步/异步的驱动函数来喂入读取到的数据,核心解析逻辑只写一次:
import struct import asyncio from typing import Generator, Callable, Awaitable def _parse_generator() -> Generator[int, bytes, dict]: # 核心解析逻辑:仅描述需要读取的字节数和数据处理规则 header = yield 2 start, length = struct.unpack("BB", header) if start != 0x68: raise ValueError("Unexpected start byte") data = yield length checksum = yield 1 # 这里加入校验和验证、数据转换等逻辑 return {"data": data, "checksum": checksum, "length": length} # 同步驱动函数 def parse_frame_sync(read_n: Callable[[int], bytes]) -> dict: gen = _parse_generator() next(gen) try: while True: required_bytes = gen.send(None) data = read_n(required_bytes) result = gen.send(data) if isinstance(result, dict): return result except StopIteration as e: return e.value # 异步驱动函数 async def parse_frame_async(read_n: Callable[[int], Awaitable[bytes]]) -> dict: gen = _parse_generator() next(gen) try: while True: required_bytes = gen.send(None) data = await read_n(required_bytes) result = gen.send(data) if isinstance(result, dict): return result except StopIteration as e: return e.value # 自动适配入口 def auto_parse_frame(read_n: Callable[[int], bytes | Awaitable[bytes]]): if asyncio.iscoroutinefunction(read_n): return parse_frame_async(read_n) else: return parse_frame_sync(read_n)
使用示例
- 同步场景:
def sync_read(n: int) -> bytes: return your_sync_stream.readexactly(n) result = auto_parse_frame(sync_read)
- 异步场景:
async def async_read(n: int) -> bytes: return await your_async_stream.readexactly(n) result = await auto_parse_frame(async_read)
方案二:类封装的IO适配
通过类把核心解析逻辑抽为静态方法,同步/异步方法仅处理IO调用差异,代码结构更清晰,扩展性更强:
import struct import asyncio from typing import Protocol, Any # 定义同步/异步读取器的协议接口 class SyncReadable(Protocol): def readexactly(self, n: int) -> bytes: ... class AsyncReadable(Protocol): async def readexactly(self, n: int) -> bytes: ... class FrameParser: @staticmethod def _core_parse(header: bytes, data: bytes, checksum: bytes) -> dict: # 纯同步的核心解析逻辑,无IO操作 start, length = struct.unpack("BB", header) if start != 0x68: raise ValueError("Unexpected start byte") # 加入校验和验证等逻辑 return {"data": data, "checksum": checksum, "length": length} @classmethod def parse_sync(cls, reader: SyncReadable) -> dict: header = reader.readexactly(2) start, length = struct.unpack("BB", header) data = reader.readexactly(length) checksum = reader.readexactly(1) return cls._core_parse(header, data, checksum) @classmethod async def parse_async(cls, reader: AsyncReadable) -> dict: header = await reader.readexactly(2) start, length = struct.unpack("BB", header) data = await reader.readexactly(length) checksum = await reader.readexactly(1) return cls._core_parse(header, data, checksum) @classmethod def auto_parse(cls, reader: Any): if not hasattr(reader, "readexactly"): raise TypeError("Reader must implement readexactly method") if asyncio.iscoroutinefunction(reader.readexactly): return cls.parse_async(reader) else: return cls.parse_sync(reader)
使用示例
- 同步场景:
result = FrameParser.auto_parse(your_sync_reader)
- 异步场景:
result = await FrameParser.auto_parse(your_async_reader)
内容的提问来源于stack exchange,提问作者S. Schlüters
相关产品推荐
相关产品推荐

