如何在FastAPI中记录原始HTTP请求与响应内容?
FastAPI 特定路由请求响应审计落地方案(低延迟适配K8s场景)
核心思路是审计逻辑完全旁路,主链路只做最少操作,绝对不把审计IO操作放到请求响应的主路径上,避免拖慢接口耗时。
1. 主链路(FastAPI层)实现要求
主链路只做2件低耗时操作,整体耗时控制在微秒级,完全不影响接口响应:
- 为特定路由打标记,只对标记路由采集原始请求、响应的字节流,不要全量路由采集引入额外开销
- 把采集到的原始数据带上trace_id、时间戳、路由信息等元数据后,直接写入本地高性能队列,立刻返回响应
不推荐直接使用FastAPI自带的
BackgroundTasks处理1MB级别的审计数据落盘:BackgroundTasks生命周期绑定到当前worker进程,若K8s触发Pod重启、滚动更新时,未执行完的后台任务会直接丢失,且大量堆积的后台任务会抢占业务逻辑的CPU内存资源。
注意:本地队列要设置最大长度阈值,队列满时直接丢弃审计数据(或降级写本地临时文件),绝对不能因为队列阻塞主请求,审计场景默认优先保障业务可用性,允许极小概率的数据丢失。
2. 审计数据落地实现(适配K8s部署)
根据你的存储需求二选一即可:
方案一:Sidecar模式(K8s原生适配,无业务侵入)
- 给Pod挂载
emptyDir类型的共享存储卷,主业务容器只需要把审计数据写入共享卷的本地文件,写本地文件的耗时极低,完全不阻塞主链路 - 额外给Pod加一个审计专用sidecar容器,只负责监听共享目录的新增文件,批量压缩后写入你需要的存储组件(Elasticsearch、对象存储、专用审计库都可以),sidecar的资源配额单独配置,就算sidecar崩溃也完全不影响主业务容器
方案二:异步消息队列模式
- 主链路把审计数据写入本地消息队列生产者的内存缓冲,由生产者客户端异步批量发往Kafka/RabbitMQ集群,发送逻辑完全不阻塞主链路
- 独立部署审计消费服务,从消息队列拉取数据后落地存储,扩容、降级完全和主业务服务解耦
3. 1MB大Body场景优化点
- 主链路不要对请求、响应Body做任何JSON解析、校验操作,直接存原始字节流即可,解析逻辑全部丢给下游的审计处理模块执行
- 审计数据落存储前做gzip压缩,1MB的JSON压缩后通常只有几十KB,大幅降低存储成本
- 所有审计逻辑的异常要单独捕获打日志,绝对不能抛出到主业务链路导致接口报错
4. FastAPI侧核心代码示例
from fastapi import FastAPI, Request, Response import threading import time from collections import deque from typing import Callable # 全局内存队列,可根据实际并发调整最大长度,满了自动丢最早的数据 AUDIT_QUEUE = deque(maxlen=1000) # 本地审计日志路径,对应K8s挂载的emptyDir共享卷 AUDIT_LOG_PATH = "/audit_shared_dir" def audit_consumer_thread(): """独立消费线程写审计数据到本地共享目录,完全不占用业务协程资源""" while True: if AUDIT_QUEUE: item = AUDIT_QUEUE.popleft() trace_id = item["trace_id"] ts = int(time.time() * 1000) # 直接写原始字节流,不做任何解析操作 with open(f"{AUDIT_LOG_PATH}/{trace_id}_{ts}.bin", "wb") as f: f.write(item["req_body"] + b"|||" + item["resp_body"]) else: time.sleep(0.01) # 服务启动时启动消费线程 threading.Thread(target=audit_consumer_thread, daemon=True).start() app = FastAPI() async def audit_middleware(request: Request, call_next: Callable): # 只处理需要审计的特定路由,可根据实际规则调整 if not request.url.path.startswith("/need_audit/"): return await call_next(request) # 采集原始请求body req_body = await request.body() # 执行业务逻辑 response = await call_next(request) # 采集原始响应body resp_body = b"" async for chunk in response.body_iterator: resp_body += chunk # 直接塞队列,不做任何IO操作 AUDIT_QUEUE.append({ "trace_id": request.headers.get("x-trace-id", "unknown"), "req_body": req_body, "resp_body": resp_body }) # 重构响应返回 return Response( content=resp_body, status_code=response.status_code, headers=dict(response.headers), media_type=response.media_type ) app.middleware("http")(audit_middleware)
内容的提问来源于stack exchange,提问作者Martín R. Niño
相关产品推荐
相关产品推荐

