From f99d587d739ee9da1c777c67e85ba548416b9a0f Mon Sep 17 00:00:00 2001 From: zhangtao <9480807882@qq.com> Date: Fri, 24 Jan 2025 20:23:17 +0800 Subject: [PATCH] =?UTF-8?q?=E5=90=8E=E7=AB=AF=E5=BC=95=E5=85=A5BackgroundT?= =?UTF-8?q?asks=E3=80=81Celert=E3=80=81websocket=E4=B8=89=E7=A7=8D?= =?UTF-8?q?=E7=B1=BB=E5=9E=8B=E7=9A=84=E6=8E=A5=E5=8F=A3=EF=BC=8C=E4=B8=BA?= =?UTF-8?q?=E5=90=8E=E6=9C=9F=E5=BA=94=E7=94=A8=E9=83=A8=E5=88=86=E5=81=9A?= =?UTF-8?q?=E9=A2=84=E7=95=99?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../v1/controllers/system/auth_controller.py | 35 +++++++++++++++++- backend/app/core/tasks.py | 24 ++++++++++++ backend/dev_sql.db | Bin 450560 -> 450560 bytes backend/requirements.txt | 3 ++ 4 files changed, 60 insertions(+), 2 deletions(-) create mode 100644 backend/app/core/tasks.py diff --git a/backend/app/api/v1/controllers/system/auth_controller.py b/backend/app/api/v1/controllers/system/auth_controller.py index 11def157..1026e484 100644 --- a/backend/app/api/v1/controllers/system/auth_controller.py +++ b/backend/app/api/v1/controllers/system/auth_controller.py @@ -1,12 +1,13 @@ # -*- coding: utf-8 -*- +import asyncio from typing import Union, Dict -from fastapi import APIRouter, Depends, Request +from fastapi import APIRouter, BackgroundTasks, Depends, Request, WebSocket from fastapi.responses import JSONResponse from sqlalchemy.ext.asyncio import AsyncSession from app.config.setting import settings -from app.common.response import ErrorResponse, SuccessResponse +from app.common.response import ErrorResponse, SuccessResponse, StreamResponse from app.api.v1.services.system.auth_service import ( LoginService, CaptchaService @@ -24,6 +25,7 @@ from app.core.dependencies import ( from app.core.router_class import OperationLogRoute from app.core.security import CustomOAuth2PasswordRequestForm from app.core.logger import logger +from app.core.tasks import background_task, long_running_task router = APIRouter(route_class=OperationLogRoute) @@ -76,3 +78,32 @@ async def logout( logger.info('退出成功') return SuccessResponse(msg='退出成功') return ErrorResponse(msg='退出失败') + +# ws://127.0.0.1:8000/api/v1/system/auth/ws +@router.websocket("/ws", name="websocket") +async def websocket_endpoint(websocket: WebSocket): + await websocket.accept() + while True: + data = await websocket.receive_text() + await websocket.send_text(f"Message text was: {data}") + + +@router.post("/celery-task", summary="模拟celery进行后台任务") +async def process_data(data: str): + # 调用 Celery 任务 + result = long_running_task.delay(data) + return SuccessResponse(msg=f"任务已提交到 Celery: {result}") + + +# 模拟流响应 +@router.get("/bg-stream", summary="模拟fastapi自带后台任务-模拟流式响应") +async def stream_response(background_tasks: BackgroundTasks): + def log_task(message): + logger.info(message) + + background_tasks.add_task(log_task, "Streaming started") + return StreamResponse( + data = background_task(), + headers={"X-Custom-Header": "Streaming-Response"}, + media_type="text/plain", + ) \ No newline at end of file diff --git a/backend/app/core/tasks.py b/backend/app/core/tasks.py new file mode 100644 index 00000000..c04a4f34 --- /dev/null +++ b/backend/app/core/tasks.py @@ -0,0 +1,24 @@ +from celery import Celery + +# 创建 Celery 实例 +app = Celery( + "worker", + broker="redis://localhost:6379/0", + backend="redis://localhost:6379/0", # 使用 Redis 存储任务结果 + include=["tasks"], # 包含任务模块 +) + +# 配置 Celery +app.conf.update( + result_expires=3600, # 任务结果过期时间(秒) + timezone="Asia/Shanghai", # 时区 +) + +@app.task +def long_running_task(data): + return f"Processed: {data}" + + +def background_task(): + for i in range(10): + yield f"执行后台任务: {i}\n" diff --git a/backend/dev_sql.db b/backend/dev_sql.db index 169ddd806884fe3338d24aaa0c916e981c9ab326..8a3e85673376b0e2c3ea0b52658b6bdf62d2c839 100644 GIT binary patch delta 1325 zcmd5)T}V?=9N*d9bTYeVRx@;Jb04=}?l|}R0~10BZH%-}=*hbeElME`SaudgL8f#XCmjB$a!Wx zpmVOJ3v{OXm;1p&hM-{E<601T-QZm#ALZ`Yun&T7J_>eN;CJ{5et@em10TZMa2O`w zW!MYD@Dw}=JM@+uXoO8R!UCwm2o^@5F#Lt#!xuh+Og-hCv(&qsKb@bQbJ(#Cj$56} zI$8&poP?hU^qwyuOAr*fh-<%s6E1R%d`GU38FCTk;d^+{ba|Gf9fTY7o;i1(DuZAd z1j@i)2EM6{U0@!ozOhR%eT{X5!*bweJ*-PPJhkmL$LzGNnQ)l*?pom_q@X42ZaCW7 zWK9kA;8d&aNTs{Nn>hW(e#qW+WpsE5Yp?C@&TQ7Wn>O@U+cUFUs|&{TlDDgvhw4R| z%t>N2MkjbBPV+J&(~1yhXcQF(Bql0pQCY(6JvOhl{XD&WZ@!dc{_WFz`o1yyz?fOu zDqhYaR-h3>v%Hr@DkH0m6yijL6mbuRL1@wygqBFPm8wqm=_SU>ahWq6h)OI$^PIxd z3LDetc%0<~S&R>`32mQV7)2G3%5fov7Z}7=c&TQD)OY||l;2(=;bJZ3Px{+IU)Jsp zRm4{4L^sQMW_I%9$-FU*w^m-N-cr`3e2a>srOkkr9MeXw=BHj5D~o%?ZwljBm6JoF z$ni3&mR!Mbsd4}RBbgOc#D-W!5{1eiW19KHh)inS_b(-1+*ishuku1jLP%t|>Xa)O cX)@N1Z^(9xS5Modshw_{dF=X5cgzl(o8|0Mo8{(Sya{;2J02N=cpCo3o@PY*rMxSmCkN1I`~zy(H)?N%2U z!?}SHOQt8@W0c$Od53W~GqWP!<>~AX7;U%vJz$LB1In7WUw+33#7x^SzhgEG