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

如何实现兼容同步/异步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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 22:45:36