ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

aiohttp 客户端 Tracing 增强:`on_response_chunk_received` 信号按块触发全量覆盖 `ClientResponse.content` 读取路径

aiohttp 客户端 Tracing 增强:`on_response_chunk_received` 信号按块触发全量覆盖 `ClientResponse.content` 读取路径 aiohttp 客户端 Tracing 增强on_response_chunk_received信号按块触发全量覆盖ClientResponse.content读取路径【免费下载链接】aiohttpAsynchronous HTTP client/server framework for asyncio and Python项目地址: https://gitcode.com/gh_mirrors/ai/aiohttp本篇技术指南围绕 aiohttpAsynchronous HTTP client/server framework for asyncio and Python中的一个重要变更展开TraceConfig.on_response_chunk_received信号从仅在ClientResponse.read()时触发一次演进为对ClientResponse.content的每一种读取方式read、readany、readchunk、readuntil、read_nowait、iter_chunked、iter_any、iter_chunks逐块触发。读者读完本文后将能理解 aiohttp 客户端追踪体系的工作方式与信号参数结构掌握利用该信号对流式响应做逐块监控、统计与审计的实战写法并从源码层面弄清单个信号回调在响应读取全路径中的触发点与调用链。变更概览一次触发到逐块触发的行为迁移本次变更记录于仓库 CHANGES/12889.bugfix.rst要点如下新行为TraceConfig.on_response_chunk_received现在会对从ClientResponse.content读取方法返回的每一个块chunk都触发一次。覆盖的读取方法read、readany、readchunk、readuntil、read_nowait、iter_chunked、iter_any、iter_chunks。旧行为此前该信号只在ClientResponse.read()时以整个响应体为一块触发一次而直接访问.content进行流式读取时追踪被完全绕过。提交者由Dreamsorcerer提交归类为 bugfix修复流式读取绕过追踪这一缺陷。这是一个典型的行为修复型变更它没有引入新 API而是让既有的追踪信号覆盖此前遗漏的读取路径使基于该信号实现的日志、指标与审计逻辑在流式场景下同样生效。为什么需要这个修复.content流式读取曾是追踪盲区aiohttp 客户端响应体有两种典型消费方式一次性读取await resp.read()、await resp.text()、await resp.json()——它们内部最终都会把整个响应体读进内存。流式读取直接操作resp.content一个StreamReader实例调用read、readany、readchunk、readuntil、iter_chunked等接口按块消费数据适用于大响应、SSEServer-Sent Events流、长连接推送等场景。修复之前on_response_chunk_received仅在ClientResponse.read()这条路径上被触发一次性携带整个 body而流式路径完全没有任何追踪事件产生。这意味着依赖该信号做响应体传输统计的应用在流式场景下统计结果恒为零想要逐块审计、限流或记录接收进度的需求无法通过现有追踪钩子实现请求端有on_request_chunk_sent请求体逐块发送信号响应端却只能整块回调追踪能力不对称。本次变更正是补上了这一缺口让响应体接收侧也具备与发送侧对称的逐块追踪能力。源码实现信号如何从StreamReader一路传到TraceConfig信号的定义与参数结构在 aiohttp/tracing.py 中TraceConfig通过 aiosignal 的Signal机制管理所有追踪钩子on_response_chunk_received对应内部信号_on_response_chunk_receivedtracing.py并对外暴露只读属性on_response_chunk_receivedtracing.py。与请求侧逐块发送信号on_request_chunk_sent相对应。该信号携带的参数类型为TraceResponseChunkReceivedParamstracing.py是一个冻结数据类包含三个字段字段类型含义methodstrHTTP 方法如GET、POSTurlyarl.URL请求的完整 URLchunkbytes本次收到的数据块信号的实际发送由Trace内部依赖持有类的send_response_chunk_received方法完成tracing.py它把method、url、chunk封装进参数对象然后调用on_response_chunk_received.send(session, trace_config_ctx, params)。回调链从响应对象挂接到流读取器触发链路的核心在 aiohttp/client_reqrep.pyClientResponse在解析响应体时把自身的_on_chunk_response_received方法挂到流读取器content的_on_chunk_received回调槽上client_reqrep.pypayload._on_chunk_received self._on_chunk_response_received_on_chunk_response_received遍历该请求关联的所有 traceself._traces对每一个 trace 调用trace.send_response_chunk_received(self.method, self.url, chunk)client_reqrep.py若任一回调抛异常会关闭连接并向上抛出。StreamReader在每次实际返回数据块前检查_on_chunk_received是否挂载若挂载则通过_fire_chunk_received触发aiohttp/streams.py。该触发运行在流自身的计时器timer上下文中意味着一个挂起hung的追踪处理器会被sock_read的读超时机制所约束不会无限拖住事件循环。各读取方法中的触发点从 aiohttp/streams.py 的源码可以看到_fire_chunk_received被下列方法在返回非空块时统一调用readuntil含其别名readline在读到分隔符、拼装出chunk后触发streams.pyread(n)从缓冲区取回chunk后触发streams.pyreadany()取回当前所有可用数据后触发streams.pyreadchunk()在 HTTP 分块边界处取回块数据后触发streams.py以及从缓冲区取整块数据后触发streams.pyread_nowait同步读取路径同样挂接回调streams.py 附近。而三个异步迭代器最终都收敛到上述底层方法iter_chunked(n)基于read(n)streams.py、iter_any基于readany()streams.py、iter_chunks基于readchunk()streams.py因此它们同样被本次变更覆盖。这也是 changelog 中列出的 8 个方法全部生效的原因——触发逻辑被下沉到了数据真正被消费的公共底层。一次触发到逐块的语义差异修复前ClientResponse.read()一次性把整个 body 交给content.read()回调只触发一次params.chunk是完整响应体直接操作.content时由于读取不经过上述挂接路径的完整回调或整体被绕过追踪事件缺失。修复后每次从流中取出并返回一个数据块都会触发一次回调params.chunk即为该次返回的块内容。因此对同一个响应回调的触发次数取决于消费方式read()一次即触发一次块为整体而iter_chunked(512)会把 4 KiB 的响应切成 8 个 512 字节的块、触发 8 次。实战示例注册逐块追踪处理器下面给出完整的可运行示例展示如何注册on_response_chunk_received处理器并观察流式响应。import asyncio import aiohttp from aiohttp import web def make_app() - web.Application: app web.Application() async def stream_handler(request: web.Request) - web.StreamResponse: resp web.StreamResponse() resp.content_length 4096 await resp.prepare(request) await resp.write(bx * 4096) return resp app.router.add_get(/, stream_handler) return app async def main() - None: # 1) 创建 TraceConfig 并注册逐块回调 trace_config aiohttp.TraceConfig() chunks: list[bytes] [] async def on_response_chunk_received( session: object, context: object, params: aiohttp.TraceResponseChunkReceivedParams, ) - None: chunks.append(params.chunk) print(f[trace] method{params.method} url{params.url} fchunk_size{len(params.chunk)}) trace_config.on_response_chunk_received.append(on_response_chunk_received) # 2) 将 trace_config 挂到 ClientSession 上 runner web.AppRunner(make_app()) await runner.setup() site web.TCPSite(runner, 127.0.0.1, 8080) await site.start() async with aiohttp.ClientSession(trace_configs[trace_config]) as session: async with session.get(http://127.0.0.1:8080/) as resp: # 3) 流式读取每个 512 字节的块都会触发一次回调 async for _ in resp.content.iter_chunked(512): pass print(ftotal chunks traced: {len(chunks)}) print(ftotal bytes traced: {sum(len(c) for c in chunks)}) await runner.cleanup() asyncio.run(main())运行这段代码控制台会打印 8 次chunk_size512的追踪记录4 KiB 响应按 512 字节切块total chunks traced为 8total bytes traced为 4096。如果把iter_chunked(512)换成await resp.read()则只打印 1 次、chunk_size4096——这直观展示了逐块触发与一次触发两种语义。常见应用场景传输量统计与审计累加params.chunk的长度得到实际经追踪路径消费的字节数可配合on_request_chunk_sent实现请求/响应双侧的对称计量。流式响应调试记录每个块到达的时间点与大小定位大响应卡顿发生在哪个块。自定义限流/放行策略在回调中基于累计字节数决定是否继续读取。指标上报将params.url与params.method作为标签逐块累加后周期上报给监控系统。注意回调是协程函数必须await内部操作且如上文所述回调执行耗时受流读取超时机制约束不要在回调中做阻塞式重活。测试验证仓库如何证明该行为本次变更配套的测试集中在 tests/test_client_session.pytest_response_chunk_received_via_contenttest_client_session.py服务端写入 4096 字节流式响应客户端通过resp.content.iter_chunked(512)消费断言每个块都被收集进chunks列表。这正是对直接.content访问不再绕过追踪这一修复点的直接回归验证。同文件中的其他相关测试覆盖了read/text/json等一次性读取路径下信号恰好触发一次assert_called_once_with风格断言以及通过on_response_chunk_received把整个响应体拼接还原test_client_session.py 附近的 gather 测试收集的字节串与响应体完全一致。tests/test_tracing.py 则断言TraceConfig冻结后on_response_chunk_received信号处于 frozen 状态保证会话复用期间信号不会被误改。总结与迁移建议对于已经在使用on_response_chunk_received的应用如果代码里通过resp.read()/resp.text()/resp.json()一次性消费响应行为基本不变仍触发一次块为整个 body可平滑升级。如果代码里使用.content流式读取且此前依赖该信号不触发的隐性行为例如以回调次数推断响应是否走流式升级后回调会按块触发需要适配。如果此前因追踪盲区而未在流式场景使用该信号现在可以直接启用无需额外改动读取代码。从源码结构看本次变更把响应体追踪下沉到StreamReader的数据消费底层aiohttp/streams.py并通过ClientResponse._on_chunk_response_received这个挂接点aiohttp/client_reqrep.py实现与追踪配置的解耦——这也意味着未来若新增其他流式读取方法只要走同一底层即可自动获得逐块追踪能力。【免费下载链接】aiohttpAsynchronous HTTP client/server framework for asyncio and Python项目地址: https://gitcode.com/gh_mirrors/ai/aiohttp创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表