FastAPI首次请求阻塞需重试问题排查求助
问题现象
搭建了一个基于FastAPI的服务,用于创建worker进程处理数据。出现异常:首次发起请求时,请求会一直挂起,无法到达API处理逻辑;取消该请求后再次发起,一切正常,worker进程能正常启动。
相关代码
创建worker的端点代码
@app.get("/connectors") async def start_process(conf: Union[Connector, None] = None): print(conf) conf_list = conf.config[0] if not conf_list: return {"message": "Empty configuration list"} log.info(f"Cassandra Host: {conf_list.cassandra_host}, Type: {type(conf_list.cassandra_host)}") globalConfig.set(conf_list.cassandra_host, conf_list.cassandra_port, conf_list.cassandra_keyspace, conf_list.kafka_bootstrap_server, conf_list.topic) # Check if the process with the same name already exists existing_process = next((p for p in processes if p["name"] == conf.name), None) if existing_process: return {"message": "Connector is not created. Duplicated error."} p = multiprocessing.Process(target=main, args=(len(processes) + 1,conf_list.connector_class,globalConfig,)) p.start() processes.append({"p": p, "name": conf.name}) return {"message": f"Started process with ID {p.pid} or {conf.name}"}
请求体示例
{ "name":"cassandra-sink-dog1", "config":[{ "connector_class":"io.inovasyon.connect.cassandra.ShardedSinkConnector", "cassandra_host":"192.168.2.165", "cassandra_port":"9042", "cassandra_keyspace":"orion_dog", "kafka_bootstrap_server":"192.168.2.165:29092", "topic":"eys.orion-dog.entities" } ] }
完整入口代码
from fastapi import FastAPI import multiprocessing from config.Config import Config from typing import List,Union from pydantic import BaseModel from functions.workerReplication import ReplicationSinkConnector from functions.workerShard import KafkaConsumerSinkConnector from functions.workerLogstash import KafkaConsumerLogstashConnector import logging import coloredlogs import uvicorn from termcolor import colored coloredlogs.install(level="INFO") log = logging.getLogger("Functions") app = FastAPI() processes = [] globalConfig = Config() class ConnectorConf(BaseModel): connector_class: str cassandra_port: str cassandra_host: str cassandra_keyspace: str kafka_bootstrap_server: str topic: str class Connector(BaseModel): name: str config: List[ConnectorConf] def main(num,classType,globalConfig): """A function that simulates a process""" log.info(colored("Worker: ","blue")+ "{}".format(num) + colored(" Running","green")) if "ShardedSinkConnector" in str(classType): log.info(colored("ShardedSinkConnector: ","green")+ colored(" started.","green")) KafkaConsumerSinkConnector().workerKafkaShard(globalConfig) if "ReplicationSinkConnector" in str(classType): log.info(colored("ReplicationSinkConnector: ","green")+ colored(" started.","green")) ReplicationSinkConnector.workerCassandraReplication(globalConfig) if "LogstashSinkConnector" in str(classType): log.info(colored("LogstashSinkConnector: ","green")+ colored(" started.","green")) KafkaConsumerLogstashConnector.workerKafkaLogstash(globalConfig,classType) @app.get("/connectors") async def start_process(conf: Union[Connector, None] = None): conf_list = conf.config[0] log.info(f"Cassandra Host: {conf_list.cassandra_host}, Type: {type(conf_list.cassandra_host)}") globalConfig.set(conf_list.cassandra_host, conf_list.cassandra_port, conf_list.cassandra_keyspace, conf_list.kafka_bootstrap_server, conf_list.topic) # Check if the process with the same name already exists existing_process = next((p for p in processes if p["name"] == conf.name), None) if existing_process: return {"message": "Connector is not created. Duplicated error."} p = multiprocessing.Process(target=main, args=(len(processes) + 1,conf_list.connector_class,globalConfig,)) p.start() processes.append({"p": p, "name": conf.name}) return {"message": f"Started process with ID {p.pid} or {conf.name}"} @app.get("/connectors/{name}") async def stop_process(name: Union[str, None] = None): for idx, p in enumerate(processes): if p["name"] == name: print(p["p"]) process_id = p["p"].pid p = p["p"] p.terminate() processes.pop(idx) log.info(f"Stopped process with ID: {process_id} - {name} Stopped") return {"message": f"Process with ID {id} not found"} # if __name__ == "__main__": # uvicorn.run(app, host="0.0.0.0", port=8000) # python -m uvicorn CassandraBackend:app --host 0.0.0.0 --port 8000 --reload
问题根源
1. GET请求携带请求体的兼容性问题
/connectors端点使用@app.get注解,但接收Connector类型的请求体参数。HTTP标准中GET请求不应该携带请求体,虽然FastAPI支持该行为,但部分HTTP客户端、服务器或代理对GET带请求体的处理逻辑存在缺陷,首次请求时可能导致请求解析阻塞,无法进入处理函数;取消后重试时,连接状态或解析逻辑变化,才正常处理请求。
2. 全局对象跨进程传递的序列化阻塞
启动子进程时直接传递globalConfig对象,如果Config类包含不可序列化的资源(如网络连接、文件句柄),首次序列化/反序列化该对象时会阻塞主线程,导致请求挂起。取消请求后部分资源释放,再次启动时避开阻塞点。
3. 缺少参数校验与异常处理
处理函数直接访问conf.config[0],若首次请求时conf为None(请求体解析失败),会触发未捕获异常,而FastAPI在首次请求时的进程初始化状态下,无法及时返回异常响应,导致请求挂起。
修复方案
1. 将GET端点改为POST
遵循REST规范,创建资源用POST请求,避免GET带请求体的兼容性问题:
@app.post("/connectors") async def start_process(conf: Connector): # 强制要求请求体,去掉可选的None # 后续逻辑不变
2. 避免传递全局对象,改用原始参数传递
在子进程中重新初始化Config对象,传递原始配置参数而非全局对象,避免序列化问题:
# 修改进程启动代码 p = multiprocessing.Process( target=main, args=( len(processes) + 1, conf_list.connector_class, (conf_list.cassandra_host, conf_list.cassandra_port, conf_list.cassandra_keyspace, conf_list.kafka_bootstrap_server, conf_list.topic) ) ) # 修改main函数 def main(num, classType, config_args): globalConfig = Config() globalConfig.set(*config_args) # 后续业务逻辑不变
3. 增加参数校验与异常捕获
确保请求参数合法,捕获处理过程中的异常,及时返回错误响应:
@app.post("/connectors") async def start_process(conf: Connector): try: if not conf.config: return {"message": "Empty configuration list"} conf_list = conf.config[0] if not conf_list: return {"message": "Invalid configuration item"} # 检查重复连接器 existing_process = next((p for p in processes if p["name"] == conf.name), None) if existing_process: return {"message": "Connector is not created. Duplicated error."} # 启动worker进程 p = multiprocessing.Process( target=main, args=(len(processes) + 1, conf_list.connector_class, (conf_list.cassandra_host, conf_list.cassandra_port, conf_list.cassandra_keyspace, conf_list.kafka_bootstrap_server, conf_list.topic)) ) p.start() processes.append({"p": p, "name": conf.name}) return {"message": f"Started process with ID {p.pid} or {conf.name}"} except Exception as e: log.error(f"Error starting process: {str(e)}") return {"message": f"Failed to start process: {str(e)}"}, 500
4. 调整服务启动方式
启用if __name__ == "__main__"包裹uvicorn启动逻辑,避免多进程环境下的重复初始化问题,同时避免使用reload模式:
if __name__ == "__main__": uvicorn.run(app, host="0.0.0.0", port=8000)
使用python your_file.py命令启动服务,而非uvicorn命令行启动,避免reload模式创建的子进程与worker进程冲突。
内容的提问来源于stack exchange,提问作者rumeysa yuk

