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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 08:15:05