根因: 1. pump() 仅在 addFiles() 调用一次,上传完成/失败后不重新触发, 导致 CONCURRENCY=3 之后的文件(4、5)永远不启动 → 无 DB 会话 → 刷新后消失 2. write_chunk 的 uploaded_chunks 是 read-modify-write,并发分片写入 后者覆盖前者 → 分片记录丢失 → complete 报 409 3. upload_chunk 是 async 但同步调 write_chunk(fsync+DB commit), 阻塞 uvicorn 事件循环 → 所有 web 请求被串行化 修复: - _shared.py: pump() 加 finally 块,上传完成/失败后都触发下一文件; 文件并发(FILE_CONCURRENCY=5)与分片并发(CHUNK_CONCURRENCY=3)分离 - upload_service.py: 按 upload_id 的进程级锁串行化 uploaded_chunks 更新, 持锁后 db.refresh 重读最新值再 append,杜绝丢失更新;complete 后清理锁 - upload_router.py: upload_chunk 的 write_chunk 调用改用 run_in_threadpool, 阻塞 I/O 移出事件循环,web 请求不再被分片写入阻塞 验证:5 文件并发上传后刷新全部可见;4 任务并发处理 4/4 成功(19.2s)
75 lines
2.6 KiB
Python
75 lines
2.6 KiB
Python
"""分片上传路由:建会话 / 查状态 / 传分片 / complete。
|
||
|
||
complete 成功后创建转写 Task 并交由 scheduler 入队。
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
from fastapi import APIRouter, Depends, Request
|
||
from starlette.concurrency import run_in_threadpool
|
||
from sqlalchemy.orm import Session
|
||
|
||
from ..database import get_db
|
||
from ..schemas.task import (
|
||
ChunkUploadResponse,
|
||
CompleteResponse,
|
||
CreateSessionRequest,
|
||
CreateSessionResponse,
|
||
SessionStatusResponse,
|
||
)
|
||
from ..services.upload_service import UploadService
|
||
|
||
router = APIRouter(prefix="/api/tasks/chunk-uploads", tags=["upload"])
|
||
|
||
|
||
def _service(db: Session = Depends(get_db)) -> UploadService:
|
||
return UploadService(db)
|
||
|
||
|
||
@router.post("", response_model=CreateSessionResponse, summary="创建分片上传会话")
|
||
def create_session(
|
||
body: CreateSessionRequest,
|
||
service: UploadService = Depends(_service),
|
||
) -> CreateSessionResponse:
|
||
return service.create_session(body)
|
||
|
||
|
||
@router.get("/{upload_id}/status", response_model=SessionStatusResponse, summary="查询会话状态(断点续传)")
|
||
def session_status(
|
||
upload_id: str,
|
||
service: UploadService = Depends(_service),
|
||
) -> SessionStatusResponse:
|
||
return service.get_status(upload_id)
|
||
|
||
|
||
@router.post("/{upload_id}/chunks/{index}", response_model=ChunkUploadResponse, summary="上传单个分片")
|
||
async def upload_chunk(
|
||
upload_id: str,
|
||
index: int,
|
||
request: Request,
|
||
service: UploadService = Depends(_service),
|
||
) -> ChunkUploadResponse:
|
||
# write_chunk 做文件 fsync + DB commit(阻塞 I/O),必须放到线程池跑,
|
||
# 否则会阻塞 uvicorn 事件循环,导致并发分片上传被串行化、web 请求卡顿。
|
||
body = await request.body()
|
||
uploaded = await run_in_threadpool(service.write_chunk, upload_id, index, body)
|
||
return ChunkUploadResponse(upload_id=upload_id, index=index, uploaded_chunks=uploaded)
|
||
|
||
|
||
@router.post("/{upload_id}/complete", response_model=CompleteResponse, summary="完成拼接并创建转写任务")
|
||
def complete_session(
|
||
upload_id: str,
|
||
service: UploadService = Depends(_service),
|
||
) -> CompleteResponse:
|
||
"""拼接分片 + 创建转写任务 + 入队管线。
|
||
|
||
controller 负责编排:service.complete 只管存储(拼接 + 建 Task),
|
||
管线触发由 controller 调用 scheduler(ffmpeg 异步 + GPU 串行调度),service 不依赖
|
||
scheduler(避免循环依赖)。
|
||
"""
|
||
resp = service.complete(upload_id)
|
||
# 仅新建任务时入队(幂等 complete 返回的也是同一 task_id,enqueue 幂等无副作用)
|
||
from ..services.scheduler import enqueue_task
|
||
enqueue_task(resp.task_id)
|
||
return resp
|