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
相关产品推荐
相关产品推荐

