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

FastAPI内部调用时中间件捕获request.body()出现阻塞问题

问题根源

FastAPI的Request.body()是一次性异步流,一旦读取后流会耗尽,后续再调用body()或者依赖请求体的代码会因等待数据而阻塞。你的中间件读取请求体后没有重置回Request对象,导致后续流程(包括内部调用其他服务时的请求转发)无法正常获取请求体。

修复方案

修改中间件,读取请求体时保存原始字节数据,读取完成后立即将原始字节重置回Request对象,确保后续流程能正常读取。

修改后的完整代码
class RequestResponseLoggingMiddleware(BaseHTTPMiddleware):
    async def dispatch(self, request: Request, call_next):
        tracer = trace.get_tracer(__name__)

        # Extract trace context from incoming request headers
        context = extract(request.headers)

        # Start a new span with the extracted context
        with tracer.start_as_current_span(f"HTTP {request.method} {request.url.path}", context=context) as span:
            start_time = time.time()

            # Extract request information
            user_id = request.headers.get("X-USER-ID", "")
            session_id = request.headers.get("X-SESSION-ID", "")
            request_id = request.headers.get("X-REQUEST-ID", "")

            # 读取原始请求体字节并保存
            raw_request_body = await request.body()
            # 处理成可读的日志格式
            try:
                request_body = raw_request_body.decode('utf-8')
            except Exception:
                request_body = str(raw_request_body)
            
            # 重置请求体,让后续的call_next或endpoint能正常读取
            await request.set_body(raw_request_body)

            # Inject trace context into outgoing request headers
            new_header = MutableHeaders(request._headers)
            inject(new_header)
            request._headers = new_header
            request.scope.update(headers=request.headers.raw)

            # Call the next middleware or endpoint
            response = await call_next(request)

            # Capture response body
            response_body = b""
            async for chunk in response.body_iterator:
                response_body += chunk

            # Create a new response object with the same content
            new_response = StreamingResponse(BytesIO(response_body), status_code=response.status_code,
                                             headers=dict(response.headers))

            try:
                response_body = response_body.decode('utf-8')
            except Exception:
                response_body = str(response_body)

            # Prepare log entry
            span_context = span.get_span_context()
            trace_id = format(span_context.trace_id, '032x')  # Convert trace_id to hex string
            span_id = format(span_context.span_id, '016x')  # Convert span_id to hex string

            parent_span_id = None
            if span.parent:
                parent_span_context = span.parent.span_context
                parent_span_id = format(parent_span_context.span_id, '016x')

            # Prepare log entry
            log_entry = {
                "timestamp": int(start_time * 1000),  # Convert to milliseconds
                "user_id": user_id,
                "session_id": session_id,
                "request_id": request_id,
                "trace_id": trace_id,
                "span_id": span_id,
                "parent_span_id": parent_span_id,
                "url": str(request.url),
                "url_params": dict(request.query_params),
                "request_body": request_body,
                "response_body": response_body,
                "response_status": response.status_code
            }

            # Log the request and response details
            logging.info(log_entry)

            return new_response
关键修改点
  • 新增raw_request_body保存原始的请求体字节数据,避免解码后无法还原成可被后续流程读取的格式
  • 调用await request.set_body(raw_request_body)重置请求体,确保后续流程(如endpoint处理请求、内部服务调用时转发请求体)能正常读取
  • 保留原有的日志格式化逻辑,仅修改请求体的读取和重置部分

这样修改后,无论是单个服务还是服务间调用,每个服务的中间件都能正常读取请求体,同时不会影响后续的请求处理流程。

内容的提问来源于stack exchange,提问作者Vinay Verma

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 18:07:45